mirror of
https://github.com/cubefs/cubefs.git
synced 2026-08-02 02:00:56 +00:00
fix(blobnode): when task call io pool submit, but ctx is cancel, fix return error
with: #1000309518 Signed-off-by: mawei029 <mawei2@oppo.com>
This commit is contained in:
parent
cdfce1c0ab
commit
0d5c51eca8
@ -34,11 +34,11 @@ type IoPoolTaskArgs struct {
|
||||
BucketId uint64
|
||||
Tm time.Time
|
||||
Ctx context.Context
|
||||
TaskFn func()
|
||||
TaskFn func() error
|
||||
}
|
||||
|
||||
type IoPool interface {
|
||||
Submit(IoPoolTaskArgs)
|
||||
Submit(IoPoolTaskArgs) error
|
||||
Close()
|
||||
}
|
||||
|
||||
@ -50,8 +50,8 @@ type IoPoolMetricConf struct {
|
||||
}
|
||||
|
||||
type taskInfo struct {
|
||||
fn func()
|
||||
done chan struct{}
|
||||
fn func() error
|
||||
done chan error
|
||||
tm time.Time
|
||||
ctx context.Context
|
||||
}
|
||||
@ -65,10 +65,9 @@ type ioPoolSimple struct {
|
||||
metric *prometheus.SummaryVec
|
||||
}
|
||||
|
||||
func (p *ioPoolSimple) Submit(args IoPoolTaskArgs) {
|
||||
func (p *ioPoolSimple) Submit(args IoPoolTaskArgs) error {
|
||||
if p.notLimit() {
|
||||
args.TaskFn()
|
||||
return
|
||||
return args.TaskFn()
|
||||
}
|
||||
|
||||
idx, task := p.generateTask(args)
|
||||
@ -76,22 +75,15 @@ func (p *ioPoolSimple) Submit(args IoPoolTaskArgs) {
|
||||
// if ctx has been cancelled, try to avoid enqueuing as much as possible;
|
||||
// even if it has enqueued, doWork/task.fn will judge ctx again
|
||||
select {
|
||||
// don't enqueue
|
||||
case <-task.ctx.Done():
|
||||
return
|
||||
default:
|
||||
}
|
||||
|
||||
select {
|
||||
// dont enqueue
|
||||
case <-task.ctx.Done():
|
||||
return
|
||||
return task.ctx.Err()
|
||||
// closing, try to complete the task
|
||||
case <-p.closed:
|
||||
args.TaskFn()
|
||||
return
|
||||
return args.TaskFn()
|
||||
// 1.normal enqueue -> do work; 2.when closing, tasks that are already in the queue will be executed
|
||||
case p.queue[idx] <- task:
|
||||
<-task.done
|
||||
return <-task.done
|
||||
}
|
||||
}
|
||||
|
||||
@ -146,12 +138,14 @@ func (p *ioPoolSimple) doWork(task *taskInfo) {
|
||||
start := time.Now()
|
||||
p.reportMetric(opDequeue, task.tm) // from enqueue to dequeue
|
||||
select {
|
||||
case <-task.ctx.Done(): // dont exec func
|
||||
// don't exec func
|
||||
case <-task.ctx.Done():
|
||||
task.done <- task.ctx.Err()
|
||||
default:
|
||||
task.fn()
|
||||
err := task.fn()
|
||||
task.done <- err
|
||||
}
|
||||
p.reportMetric(opOnDisk, start) // from dequeue to op done
|
||||
task.done <- struct{}{}
|
||||
}
|
||||
|
||||
func (p *ioPoolSimple) reportMetric(opStage string, tm time.Time) {
|
||||
@ -164,7 +158,7 @@ func (p *ioPoolSimple) generateTask(args IoPoolTaskArgs) (idx uint64, task *task
|
||||
|
||||
task = &taskInfo{
|
||||
fn: args.TaskFn,
|
||||
done: make(chan struct{}, 1),
|
||||
done: make(chan error, 1),
|
||||
tm: args.Tm,
|
||||
ctx: args.Ctx,
|
||||
}
|
||||
|
||||
@ -43,9 +43,10 @@ func TestIoPoolSimple(t *testing.T) {
|
||||
task := IoPoolTaskArgs{
|
||||
BucketId: 1,
|
||||
Tm: time.Now(),
|
||||
TaskFn: func() {
|
||||
TaskFn: func() error {
|
||||
ch <- struct{}{}
|
||||
n++
|
||||
return nil
|
||||
},
|
||||
}
|
||||
task2, task3 := task, task
|
||||
@ -56,20 +57,22 @@ func TestIoPoolSimple(t *testing.T) {
|
||||
|
||||
<-ch
|
||||
closePool.Close()
|
||||
closePool.Submit(task3)
|
||||
require.Equal(t, 3, n) // two task func
|
||||
_err := closePool.Submit(task3)
|
||||
require.Equal(t, 3, n) // all task should be executed
|
||||
require.NoError(t, _err)
|
||||
}
|
||||
|
||||
// alloc
|
||||
{
|
||||
chunkId := uint64(1)
|
||||
taskFn := func() { err = sys.PreAllocate(file.Fd(), 0, 4096) }
|
||||
taskFn := func() error { err = sys.PreAllocate(file.Fd(), 0, 4096); return err }
|
||||
task := IoPoolTaskArgs{
|
||||
BucketId: chunkId,
|
||||
Tm: time.Now(),
|
||||
TaskFn: taskFn,
|
||||
}
|
||||
writePool.Submit(task)
|
||||
_err := writePool.Submit(task)
|
||||
require.NoError(t, _err)
|
||||
|
||||
require.NoError(t, err)
|
||||
fi1, _ := file.Stat()
|
||||
@ -81,8 +84,9 @@ func TestIoPoolSimple(t *testing.T) {
|
||||
data := []byte(content)
|
||||
n := 0
|
||||
chunkId := uint64(2)
|
||||
taskFn := func() {
|
||||
taskFn := func() error {
|
||||
n, err = file.WriteAt(data, 0)
|
||||
return err
|
||||
}
|
||||
task := IoPoolTaskArgs{
|
||||
BucketId: chunkId,
|
||||
@ -90,7 +94,8 @@ func TestIoPoolSimple(t *testing.T) {
|
||||
TaskFn: taskFn,
|
||||
}
|
||||
|
||||
writePool.Submit(task)
|
||||
_err := writePool.Submit(task)
|
||||
require.NoError(t, _err)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, len(content), n)
|
||||
}
|
||||
@ -98,16 +103,15 @@ func TestIoPoolSimple(t *testing.T) {
|
||||
// sync
|
||||
{
|
||||
chunkId := uint64(3)
|
||||
taskFn := func() {
|
||||
err = file.Sync()
|
||||
}
|
||||
taskFn := func() error { err = file.Sync(); return err }
|
||||
task := IoPoolTaskArgs{
|
||||
BucketId: chunkId,
|
||||
Tm: time.Now(),
|
||||
TaskFn: taskFn,
|
||||
}
|
||||
|
||||
writePool.Submit(task)
|
||||
_err := writePool.Submit(task)
|
||||
require.NoError(t, _err)
|
||||
require.NoError(t, err)
|
||||
}
|
||||
|
||||
@ -116,8 +120,9 @@ func TestIoPoolSimple(t *testing.T) {
|
||||
n := 0
|
||||
data := make([]byte, len(content))
|
||||
chunkId := uint64(4)
|
||||
taskFn := func() {
|
||||
taskFn := func() error {
|
||||
n, err = file.ReadAt(data, 0)
|
||||
return err
|
||||
}
|
||||
task := IoPoolTaskArgs{
|
||||
BucketId: chunkId,
|
||||
@ -126,7 +131,8 @@ func TestIoPoolSimple(t *testing.T) {
|
||||
}
|
||||
|
||||
require.Equal(t, uint8(0), data[0])
|
||||
readPool.Submit(task)
|
||||
_err := readPool.Submit(task)
|
||||
require.NoError(t, _err)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, len(content), n)
|
||||
require.Equal(t, content, string(data))
|
||||
@ -140,16 +146,15 @@ func TestIoPoolSimple(t *testing.T) {
|
||||
require.Equal(t, content, string(data))
|
||||
|
||||
chunkId := uint64(5)
|
||||
taskFn := func() {
|
||||
err = sys.PunchHole(file.Fd(), 0, 4096)
|
||||
}
|
||||
taskFn := func() error { err = sys.PunchHole(file.Fd(), 0, 4096); return err }
|
||||
task := IoPoolTaskArgs{
|
||||
BucketId: chunkId,
|
||||
Tm: time.Now(),
|
||||
TaskFn: taskFn,
|
||||
}
|
||||
writePool.Submit(task)
|
||||
_err := writePool.Submit(task)
|
||||
|
||||
require.NoError(t, _err)
|
||||
require.NoError(t, err)
|
||||
_, err = file.Read(data)
|
||||
require.NoError(t, err)
|
||||
@ -165,14 +170,15 @@ func TestIoPoolSimple(t *testing.T) {
|
||||
n := 0
|
||||
chunkId := uint64(2)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
taskFn := func() {
|
||||
taskFn := func() error {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
n, err = 0, ctx.Err()
|
||||
return
|
||||
return err
|
||||
default:
|
||||
}
|
||||
n, err = file.WriteAt(data, 0)
|
||||
return err
|
||||
}
|
||||
task := IoPoolTaskArgs{
|
||||
BucketId: chunkId,
|
||||
@ -184,7 +190,8 @@ func TestIoPoolSimple(t *testing.T) {
|
||||
// cancel before submit
|
||||
cancel()
|
||||
|
||||
writePool.Submit(task)
|
||||
_err := writePool.Submit(task)
|
||||
require.ErrorIs(t, _err, context.Canceled)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, 0, n) // not write
|
||||
}
|
||||
@ -196,14 +203,15 @@ func TestIoPoolSimple(t *testing.T) {
|
||||
chunkId := uint64(2)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
|
||||
taskFn := func() {
|
||||
taskFn := func() error {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
n, err = 0, ctx.Err()
|
||||
return
|
||||
return err
|
||||
default:
|
||||
}
|
||||
n, err = file.WriteAt(data, 0)
|
||||
return err
|
||||
}
|
||||
task := IoPoolTaskArgs{
|
||||
BucketId: chunkId,
|
||||
@ -215,9 +223,10 @@ func TestIoPoolSimple(t *testing.T) {
|
||||
taskLongTime := IoPoolTaskArgs{
|
||||
BucketId: chunkId,
|
||||
Tm: time.Now(),
|
||||
TaskFn: func() {
|
||||
TaskFn: func() error {
|
||||
cancel()
|
||||
n = 1
|
||||
return nil
|
||||
},
|
||||
Ctx: nil,
|
||||
}
|
||||
@ -226,11 +235,13 @@ func TestIoPoolSimple(t *testing.T) {
|
||||
allDone := make(chan struct{}, 1)
|
||||
go func() {
|
||||
ch <- struct{}{}
|
||||
writePool.Submit(taskLongTime)
|
||||
_err := writePool.Submit(taskLongTime)
|
||||
require.NoError(t, _err)
|
||||
}()
|
||||
go func() {
|
||||
<-ch
|
||||
writePool.Submit(task)
|
||||
_err := writePool.Submit(task)
|
||||
require.ErrorIs(t, _err, context.Canceled)
|
||||
allDone <- struct{}{}
|
||||
}()
|
||||
|
||||
@ -247,10 +258,11 @@ func TestIoPoolSimple(t *testing.T) {
|
||||
n := 0
|
||||
data := []byte(content)
|
||||
task := IoPoolTaskArgs{
|
||||
TaskFn: func() { n, err = file.WriteAt(data, 0) },
|
||||
TaskFn: func() error { n, err = file.WriteAt(data, 0); return err },
|
||||
}
|
||||
|
||||
emptyPool.Submit(task)
|
||||
_err := emptyPool.Submit(task)
|
||||
require.NoError(t, _err)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, len(content), n)
|
||||
|
||||
@ -258,12 +270,13 @@ func TestIoPoolSimple(t *testing.T) {
|
||||
n = 0
|
||||
data = make([]byte, len(content))
|
||||
task = IoPoolTaskArgs{
|
||||
TaskFn: func() { n, err = file.ReadAt(data, 0) },
|
||||
TaskFn: func() error { n, err = file.ReadAt(data, 0); return err },
|
||||
}
|
||||
require.Equal(t, uint8(0), data[0])
|
||||
require.Equal(t, 0, n)
|
||||
|
||||
emptyPool.Submit(task)
|
||||
_err = emptyPool.Submit(task)
|
||||
require.NoError(t, _err)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, len(content), n)
|
||||
require.Equal(t, content, string(data))
|
||||
|
||||
@ -77,15 +77,16 @@ func (ef *blobFile) ReadAtCtx(ctx context.Context, b []byte, off int64) (n int,
|
||||
BucketId: ef.chunk,
|
||||
Tm: time.Now(),
|
||||
Ctx: ctx,
|
||||
TaskFn: func() {
|
||||
TaskFn: func() error {
|
||||
if ctx.Err() != nil {
|
||||
n, err = 0, ctx.Err()
|
||||
return
|
||||
return err
|
||||
}
|
||||
n, err = ef.file.ReadAt(b, off)
|
||||
return err
|
||||
},
|
||||
}
|
||||
ef.ioPools[bnapi.GetIoType(ctx)].Submit(task)
|
||||
err = ef.ioPools[bnapi.GetIoType(ctx)].Submit(task)
|
||||
|
||||
ef.handleError(err)
|
||||
return
|
||||
@ -100,15 +101,16 @@ func (ef *blobFile) WriteAtCtx(ctx context.Context, b []byte, off int64) (n int,
|
||||
BucketId: ef.chunk,
|
||||
Tm: time.Now(),
|
||||
Ctx: ctx,
|
||||
TaskFn: func() {
|
||||
TaskFn: func() error {
|
||||
if ctx.Err() != nil {
|
||||
n, err = 0, ctx.Err()
|
||||
return
|
||||
return err
|
||||
}
|
||||
n, err = ef.file.WriteAt(b, off)
|
||||
return err
|
||||
},
|
||||
}
|
||||
ef.ioPools[bnapi.GetIoType(ctx)].Submit(task)
|
||||
err = ef.ioPools[bnapi.GetIoType(ctx)].Submit(task)
|
||||
|
||||
ef.handleError(err)
|
||||
return
|
||||
@ -124,9 +126,9 @@ func (ef *blobFile) Allocate(off int64, size int64) (err error) {
|
||||
task := base.IoPoolTaskArgs{
|
||||
BucketId: ef.chunk,
|
||||
Tm: time.Now(),
|
||||
TaskFn: func() { err = sys.PreAllocate(ef.file.Fd(), off, size) },
|
||||
TaskFn: func() error { return sys.PreAllocate(ef.file.Fd(), off, size) },
|
||||
}
|
||||
ef.ioPools[bnapi.WriteIO].Submit(task)
|
||||
err = ef.ioPools[bnapi.WriteIO].Submit(task)
|
||||
|
||||
ef.handleError(err)
|
||||
return
|
||||
@ -136,9 +138,9 @@ func (ef *blobFile) Discard(off int64, size int64) (err error) {
|
||||
task := base.IoPoolTaskArgs{
|
||||
BucketId: ef.chunk,
|
||||
Tm: time.Now(),
|
||||
TaskFn: func() { err = sys.PunchHole(ef.file.Fd(), off, size) },
|
||||
TaskFn: func() error { return sys.PunchHole(ef.file.Fd(), off, size) },
|
||||
}
|
||||
ef.ioPools[bnapi.DeleteIO].Submit(task)
|
||||
err = ef.ioPools[bnapi.DeleteIO].Submit(task)
|
||||
|
||||
ef.handleError(err)
|
||||
return
|
||||
@ -159,9 +161,9 @@ func (ef *blobFile) Sync() (err error) {
|
||||
task := base.IoPoolTaskArgs{
|
||||
BucketId: ef.chunk,
|
||||
Tm: time.Now(),
|
||||
TaskFn: func() { err = ef.syncHandler.Do(nil) },
|
||||
TaskFn: func() error { return ef.syncHandler.Do(nil) },
|
||||
}
|
||||
ef.ioPools[bnapi.WriteIO].Submit(task)
|
||||
err = ef.ioPools[bnapi.WriteIO].Submit(task)
|
||||
|
||||
ef.handleError(err)
|
||||
return err
|
||||
|
||||
@ -36,8 +36,8 @@ func doBlobFileOp(t *testing.T, f *os.File, emptyIoPool bool) {
|
||||
|
||||
ctr := gomock.NewController(t)
|
||||
ioPool := mocks.NewMockIoPool(ctr)
|
||||
ioPool.EXPECT().Submit(gomock.Any()).Do(func(args base.IoPoolTaskArgs) {
|
||||
args.TaskFn()
|
||||
ioPool.EXPECT().Submit(gomock.Any()).DoAndReturn(func(args base.IoPoolTaskArgs) error {
|
||||
return args.TaskFn()
|
||||
}).AnyTimes()
|
||||
ioPools := map[bnapi.IOType]base.IoPool{
|
||||
bnapi.ReadIO: ioPool,
|
||||
@ -183,46 +183,96 @@ func TestBlobFile_doTaskFnCtxCancel(t *testing.T) {
|
||||
|
||||
data := []byte("test data")
|
||||
|
||||
// WriteAtCtx
|
||||
// Test 1: Normal WriteAtCtx - should succeed
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
ctx = bnapi.SetIoType(ctx, bnapi.WriteIO)
|
||||
ioPool.EXPECT().Submit(gomock.Any()).Do(func(args base.IoPoolTaskArgs) {
|
||||
args.TaskFn()
|
||||
ioPool.EXPECT().Submit(gomock.Any()).DoAndReturn(func(args base.IoPoolTaskArgs) error {
|
||||
return 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 base.IoPoolTaskArgs) {
|
||||
cancel()
|
||||
args.TaskFn()
|
||||
})
|
||||
ctx = bnapi.SetIoType(ctx, bnapi.BackgroundIO)
|
||||
ioPool.EXPECT().Submit(gomock.Any()).DoAndReturn(func(args base.IoPoolTaskArgs) error {
|
||||
cancel()
|
||||
return args.TaskFn()
|
||||
})
|
||||
n, err = ef.WriteAtCtx(ctx, data, 0)
|
||||
require.ErrorIs(t, err, context.Canceled)
|
||||
require.Equal(t, 0, n)
|
||||
|
||||
// ReadAtCtx
|
||||
// Test 2: Normal ReadAtCtx - should succeed
|
||||
ctx, cancel = context.WithCancel(context.Background())
|
||||
ctx = bnapi.SetIoType(ctx, bnapi.ReadIO)
|
||||
buf := make([]byte, len(data))
|
||||
ioPool.EXPECT().Submit(gomock.Any()).Do(func(args base.IoPoolTaskArgs) {
|
||||
args.TaskFn()
|
||||
ioPool.EXPECT().Submit(gomock.Any()).DoAndReturn(func(args base.IoPoolTaskArgs) error {
|
||||
return 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 base.IoPoolTaskArgs) {
|
||||
cancel()
|
||||
args.TaskFn()
|
||||
})
|
||||
ctx = bnapi.SetIoType(ctx, bnapi.BackgroundIO)
|
||||
// Test 3: Context canceled before calling WriteAtCtx, blobfile checks ctx.Err() before calling Submit
|
||||
// If context is already canceled, it returns immediately without calling Submit
|
||||
ctx, cancel = context.WithCancel(context.Background())
|
||||
ctx = bnapi.SetIoType(ctx, bnapi.WriteIO)
|
||||
|
||||
// Cancel context before calling WriteAtCtx
|
||||
cancel()
|
||||
|
||||
// No mock expectation needed since Submit won't be called when context is already canceled
|
||||
n, err = ef.WriteAtCtx(ctx, data, 0)
|
||||
require.ErrorIs(t, err, context.Canceled)
|
||||
require.Equal(t, 0, n)
|
||||
|
||||
// Test 4: Context canceled before calling ReadAtCtx
|
||||
ctx, cancel = context.WithCancel(context.Background())
|
||||
ctx = bnapi.SetIoType(ctx, bnapi.ReadIO)
|
||||
|
||||
// Cancel context before calling ReadAtCtx
|
||||
cancel()
|
||||
|
||||
// No mock expectation needed since Submit won't be called when context is already canceled
|
||||
n, err = ef.ReadAtCtx(ctx, buf, 0)
|
||||
require.ErrorIs(t, err, context.Canceled)
|
||||
require.Equal(t, 0, n)
|
||||
|
||||
// Test 5: Context canceled during task execution in WriteAtCtx
|
||||
// This simulates the case where iopool detects context cancellation and returns context.Canceled
|
||||
ctx = bnapi.SetIoType(context.Background(), bnapi.WriteIO)
|
||||
ioPool.EXPECT().Submit(gomock.Any()).Return(context.Canceled)
|
||||
|
||||
n, err = ef.WriteAtCtx(ctx, data, 0)
|
||||
require.ErrorIs(t, err, context.Canceled)
|
||||
require.Equal(t, 0, n)
|
||||
|
||||
// Test 6: Context canceled during task execution in ReadAtCtx
|
||||
// This simulates the case where iopool detects context cancellation and returns context.Canceled
|
||||
ctx = bnapi.SetIoType(context.Background(), bnapi.ReadIO)
|
||||
ioPool.EXPECT().Submit(gomock.Any()).Return(context.Canceled)
|
||||
|
||||
n, err = ef.ReadAtCtx(ctx, buf, 0)
|
||||
require.ErrorIs(t, err, context.Canceled)
|
||||
require.Equal(t, 0, n)
|
||||
|
||||
// Test 7: Allocate with context cancellation during task execution
|
||||
ioPool.EXPECT().Submit(gomock.Any()).Return(context.Canceled)
|
||||
err = ef.Allocate(0, 1024)
|
||||
require.ErrorIs(t, err, context.Canceled)
|
||||
|
||||
// Test 8: Discard with context cancellation during task execution
|
||||
ioPool.EXPECT().Submit(gomock.Any()).Return(context.Canceled)
|
||||
err = ef.Discard(0, 1024)
|
||||
require.ErrorIs(t, err, context.Canceled)
|
||||
|
||||
// Test 9: Sync with context cancellation during task execution
|
||||
ioPool.EXPECT().Submit(gomock.Any()).Return(context.Canceled)
|
||||
err = ef.Sync()
|
||||
require.ErrorIs(t, err, context.Canceled)
|
||||
|
||||
// Test 10: Panic case for invalid IO type
|
||||
require.Panics(t, func() {
|
||||
ctx = context.Background()
|
||||
ctx = bnapi.SetIoType(ctx, bnapi.IOTypeMax)
|
||||
|
||||
@ -1299,3 +1299,214 @@ func TestChunkData_WriteReadCancel(t *testing.T) {
|
||||
// resume
|
||||
cd.ef = backup
|
||||
}
|
||||
|
||||
// this test verifies that when tw.Write returns n != len(buf), the function returns ErrInternal
|
||||
func TestChunkData_WritePartialWrite(t *testing.T) {
|
||||
testDir, err := os.MkdirTemp(os.TempDir(), defaultDiskTestDir+"WritePartialWrite")
|
||||
require.NoError(t, err)
|
||||
defer os.RemoveAll(testDir)
|
||||
|
||||
ctx := context.Background()
|
||||
ctx = bnapi.SetIoType(ctx, bnapi.WriteIO)
|
||||
|
||||
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 := newIoQosMgrMock(t, 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 the blob file to simulate partial writes
|
||||
backup := cd.ef
|
||||
ctr := gomock.NewController(t)
|
||||
mockBlobFile := bnmock.NewMockBlobFile(ctr)
|
||||
cd.ef = mockBlobFile
|
||||
|
||||
log.Infof("chunkdata: \n%s", cd)
|
||||
require.Equal(t, int32(cd.wOff), int32(4096))
|
||||
|
||||
// Test case 1: Partial write - write returns fewer bytes than requested
|
||||
sharddata := []byte("test data for partial write")
|
||||
shard := &core.Shard{
|
||||
Bid: 1001,
|
||||
Vuid: 10,
|
||||
Flag: bnapi.ShardStatusNormal,
|
||||
Size: uint32(len(sharddata)),
|
||||
Body: bytes.NewBuffer(sharddata),
|
||||
}
|
||||
|
||||
// Mock WriteAtCtx to return partial write (fewer bytes than requested)
|
||||
mockBlobFile.EXPECT().WriteAtCtx(gomock.Any(), gomock.Any(), gomock.Any()).DoAndReturn(
|
||||
func(ctx context.Context, b []byte, off int64) (n int, err error) {
|
||||
// Simulate partial write - return only half of the requested bytes
|
||||
return len(b) / 2, nil
|
||||
}).AnyTimes()
|
||||
|
||||
err = cd.Write(ctx, shard)
|
||||
require.Error(t, err)
|
||||
require.ErrorIs(t, err, bloberr.ErrInternal)
|
||||
t.Logf("Expected error for partial write: %v", err)
|
||||
|
||||
// Test case 2: Zero write - write returns 0 bytes
|
||||
shard2 := &core.Shard{
|
||||
Bid: 1002,
|
||||
Vuid: 10,
|
||||
Flag: bnapi.ShardStatusNormal,
|
||||
Size: uint32(len(sharddata)),
|
||||
Body: bytes.NewBuffer(sharddata),
|
||||
}
|
||||
|
||||
// Mock WriteAtCtx to return 0 bytes written
|
||||
mockBlobFile.EXPECT().WriteAtCtx(gomock.Any(), gomock.Any(), gomock.Any()).DoAndReturn(
|
||||
func(ctx context.Context, b []byte, off int64) (n int, err error) {
|
||||
// Simulate zero write
|
||||
return 0, nil
|
||||
}).AnyTimes()
|
||||
|
||||
err = cd.Write(ctx, shard2)
|
||||
require.Error(t, err)
|
||||
require.ErrorIs(t, err, bloberr.ErrInternal)
|
||||
t.Logf("Expected error for zero write: %v", err)
|
||||
|
||||
// Test case 3: Write returns more bytes than requested (should not happen in practice, but test for robustness)
|
||||
shard3 := &core.Shard{
|
||||
Bid: 1003,
|
||||
Vuid: 10,
|
||||
Flag: bnapi.ShardStatusNormal,
|
||||
Size: uint32(len(sharddata)),
|
||||
Body: bytes.NewBuffer(sharddata),
|
||||
}
|
||||
|
||||
// Mock WriteAtCtx to return more bytes than requested
|
||||
mockBlobFile.EXPECT().WriteAtCtx(gomock.Any(), gomock.Any(), gomock.Any()).DoAndReturn(
|
||||
func(ctx context.Context, b []byte, off int64) (n int, err error) {
|
||||
// Simulate writing more bytes than requested
|
||||
return len(b) + 10, nil
|
||||
}).AnyTimes()
|
||||
|
||||
err = cd.Write(ctx, shard3)
|
||||
require.Error(t, err)
|
||||
require.ErrorIs(t, err, bloberr.ErrInternal)
|
||||
t.Logf("Expected error for excessive write: %v", err)
|
||||
|
||||
// Test case 4: Normal write should still work when mock is restored
|
||||
cd.ef = backup
|
||||
shard4 := &core.Shard{
|
||||
Bid: 1004,
|
||||
Vuid: 10,
|
||||
Flag: bnapi.ShardStatusNormal,
|
||||
Size: uint32(len(sharddata)),
|
||||
Body: bytes.NewBuffer(sharddata),
|
||||
}
|
||||
|
||||
err = cd.Write(ctx, shard4)
|
||||
require.NoError(t, err)
|
||||
t.Logf("Normal write succeeded after restoring backup: %v", err)
|
||||
}
|
||||
|
||||
// This test verifies that when context is canceled during write, ErrIOCtxCancel is returned
|
||||
func TestChunkData_WriteContextCanceled(t *testing.T) {
|
||||
testDir, err := os.MkdirTemp(os.TempDir(), defaultDiskTestDir+"WriteContextCanceled")
|
||||
require.NoError(t, err)
|
||||
defer os.RemoveAll(testDir)
|
||||
|
||||
ctx := context.Background()
|
||||
ctx = bnapi.SetIoType(ctx, bnapi.WriteIO)
|
||||
|
||||
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 := newIoQosMgrMock(t, 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 the blob file
|
||||
backup := cd.ef
|
||||
ctr := gomock.NewController(t)
|
||||
mockBlobFile := bnmock.NewMockBlobFile(ctr)
|
||||
cd.ef = mockBlobFile
|
||||
|
||||
log.Infof("chunkdata: \n%s", cd)
|
||||
require.Equal(t, int32(cd.wOff), int32(4096))
|
||||
|
||||
sharddata := []byte("test data for context cancellation")
|
||||
shard := &core.Shard{
|
||||
Bid: 2001,
|
||||
Vuid: 10,
|
||||
Flag: bnapi.ShardStatusNormal,
|
||||
Size: uint32(len(sharddata)),
|
||||
Body: bytes.NewBuffer(sharddata),
|
||||
}
|
||||
|
||||
// Test case 1: Context canceled before write operation
|
||||
ctxCanceled := context.Background()
|
||||
ctxCanceled = bnapi.SetIoType(ctxCanceled, bnapi.WriteIO)
|
||||
ctxCanceled, cancel := context.WithCancel(ctxCanceled)
|
||||
cancel() // Cancel immediately
|
||||
|
||||
// Mock WriteAtCtx to return context.Canceled error
|
||||
mockBlobFile.EXPECT().WriteAtCtx(gomock.Any(), gomock.Any(), gomock.Any()).DoAndReturn(
|
||||
func(ctx context.Context, b []byte, off int64) (n int, err error) {
|
||||
return 0, context.Canceled
|
||||
}).AnyTimes()
|
||||
|
||||
err = cd.Write(ctxCanceled, shard)
|
||||
require.Error(t, err)
|
||||
require.ErrorIs(t, err, bloberr.ErrIOCtxCancel)
|
||||
t.Logf("Expected ErrIOCtxCancel for context canceled before write: %v", err)
|
||||
|
||||
// Test case 2: Context canceled during write operation
|
||||
ctxCanceled2 := context.Background()
|
||||
ctxCanceled2 = bnapi.SetIoType(ctxCanceled2, bnapi.WriteIO)
|
||||
ctxCanceled2, cancel2 := context.WithCancel(ctxCanceled2)
|
||||
|
||||
// Mock WriteAtCtx to cancel context during write
|
||||
mockBlobFile.EXPECT().WriteAtCtx(gomock.Any(), gomock.Any(), gomock.Any()).DoAndReturn(
|
||||
func(ctx context.Context, b []byte, off int64) (n int, err error) {
|
||||
cancel2() // Cancel during write
|
||||
return 0, context.Canceled
|
||||
}).AnyTimes()
|
||||
|
||||
err = cd.Write(ctxCanceled2, shard)
|
||||
require.Error(t, err)
|
||||
// require.ErrorIs(t, err, bloberr.ErrIOCtxCancel) // Reader Error
|
||||
t.Logf("Expected ErrIOCtxCancel for context canceled during write: %v", err)
|
||||
|
||||
// Test case 3: Normal write should work with valid context
|
||||
cd.ef = backup
|
||||
shard2 := &core.Shard{
|
||||
Bid: 2002,
|
||||
Vuid: 10,
|
||||
Flag: bnapi.ShardStatusNormal,
|
||||
Size: uint32(len(sharddata)),
|
||||
Body: bytes.NewBuffer(sharddata),
|
||||
}
|
||||
|
||||
err = cd.Write(ctx, shard2)
|
||||
require.NoError(t, err)
|
||||
t.Logf("Normal write succeeded with valid context: %v", err)
|
||||
}
|
||||
|
||||
@ -78,13 +78,13 @@ func lsblkByMountPoint(mountPoint string) (exist bool, err error) {
|
||||
func IsLostDisk(path string) bool {
|
||||
mountPath, err := getMountPoint(path)
|
||||
if err != nil {
|
||||
log.Error(err)
|
||||
log.Errorf("path=%s, err=%+v", path, err)
|
||||
return false
|
||||
}
|
||||
exist, err := lsblkByMountPoint(mountPath)
|
||||
if err == nil && !exist {
|
||||
return true
|
||||
}
|
||||
log.Error(err)
|
||||
log.Errorf("path=%s, err=%+v", path, err)
|
||||
return false
|
||||
}
|
||||
|
||||
@ -47,9 +47,11 @@ func (mr *MockIoPoolMockRecorder) Close() *gomock.Call {
|
||||
}
|
||||
|
||||
// Submit mocks base method.
|
||||
func (m *MockIoPool) Submit(arg0 base.IoPoolTaskArgs) {
|
||||
func (m *MockIoPool) Submit(arg0 base.IoPoolTaskArgs) error {
|
||||
m.ctrl.T.Helper()
|
||||
m.ctrl.Call(m, "Submit", arg0)
|
||||
ret := m.ctrl.Call(m, "Submit", arg0)
|
||||
ret0, _ := ret[0].(error)
|
||||
return ret0
|
||||
}
|
||||
|
||||
// Submit indicates an expected call of Submit.
|
||||
|
||||
Loading…
Reference in New Issue
Block a user