mirror of
https://github.com/OpenListTeam/OpenList.git
synced 2026-10-10 04:53:09 +08:00
b6db83ed5e
* add LinearMemory * replace mmap with LinearMemory * remove unused code * add GuardedMemory; add `min_free_memoryMB` conf * add HybridCache and StreamBuffer * log * rename SizedReadWriterAt to Section * 重构 FileStream,改用 HybridCache * 重构 HybridCache,更新方法名并添加回滚功能;优化请求和流处理逻辑 * 重构 StreamSectionReader 接口,使用HybridCache * 在 NewGuardedMemory 函数中添加了对 LinearMemory 的最终化处理,以确保内存释放 * . * 优化检查逻辑 * 重命名 * 改进、重命名 * 修复 * 添加测试 * 移除HybridCacheReader并引入DynamicReadAtSeeker * 重构缓存读取逻辑,简化代码并引入ReadFromN方法 * 优化缓存配置注释并修复下载器部分大小限制逻辑 * 优化中断逻辑 * 优化下载器代码 * 修复bug * HybridCache添加多文件缓存模式 * 优化下载器并发,添加测试 * 修复bug * 重命名+注释 * . * fix(net): always cleanup downloader on interrupt * fix(net): guard chunk enqueue with context cancel * fix(net): update interrupt logic and add download interrupt test * fix(test): update concurrency limit in high concurrency test * refactor(buffer): simplify ReadAt logic * refactor(config): update memory configuration logic * . --------- Co-authored-by: Suyunmeng <Susus0175@proton.me>
98 lines
2.1 KiB
Go
98 lines
2.1 KiB
Go
package buffer
|
|
|
|
import (
|
|
"errors"
|
|
"io"
|
|
)
|
|
|
|
type WriteAtSeekerProvider interface{ GetWriteAtSeeker() WriteAtSeeker }
|
|
|
|
func WriteAtSeekerOf(b Block) WriteAtSeeker {
|
|
if p, ok := b.(WriteAtSeekerProvider); ok {
|
|
return p.GetWriteAtSeeker()
|
|
}
|
|
return io.NewOffsetWriter(b, 0)
|
|
}
|
|
|
|
type ReadAtSeekerProvider interface{ GetReadAtSeeker() ReadAtSeeker }
|
|
|
|
// 将一个Block包装为ReadAtSeeker。
|
|
// 固定大小:当前Block的Size()。
|
|
func ReadAtSeekerOf(b Block) ReadAtSeeker {
|
|
if p, ok := b.(ReadAtSeekerProvider); ok {
|
|
return p.GetReadAtSeeker()
|
|
}
|
|
return io.NewSectionReader(b, 0, b.Size())
|
|
}
|
|
|
|
type blockAdapter struct {
|
|
WriteAtSeeker
|
|
SizedReadAtSeeker
|
|
}
|
|
|
|
func (b *blockAdapter) GetWriteAtSeeker() WriteAtSeeker {
|
|
return b.WriteAtSeeker
|
|
}
|
|
|
|
func (b *blockAdapter) GetReadAtSeeker() ReadAtSeeker {
|
|
return b.SizedReadAtSeeker
|
|
}
|
|
func NewBlockAdapter(w WriteAtSeeker, r SizedReadAtSeeker) Block {
|
|
return &blockAdapter{
|
|
WriteAtSeeker: w,
|
|
SizedReadAtSeeker: r,
|
|
}
|
|
}
|
|
|
|
var _ Block = (*blockAdapter)(nil)
|
|
|
|
// 将一个Block包装为ReadAtSeeker。
|
|
// 动态大小:Size() 是动态跟随底层 Block。
|
|
type DynamicReadAtSeeker struct {
|
|
block Block
|
|
offset int64
|
|
}
|
|
|
|
func (r *DynamicReadAtSeeker) ReadAt(p []byte, off int64) (n int, err error) {
|
|
return r.block.ReadAt(p, off)
|
|
}
|
|
|
|
func (r *DynamicReadAtSeeker) Read(p []byte) (n int, err error) {
|
|
n, err = r.block.ReadAt(p, r.offset)
|
|
if n > 0 {
|
|
r.offset += int64(n)
|
|
}
|
|
return n, err
|
|
}
|
|
|
|
func (r *DynamicReadAtSeeker) Size() int64 {
|
|
return r.block.Size()
|
|
}
|
|
|
|
func (r *DynamicReadAtSeeker) Seek(offset int64, whence int) (int64, error) {
|
|
switch whence {
|
|
case io.SeekStart:
|
|
case io.SeekCurrent:
|
|
if offset == 0 {
|
|
return r.offset, nil
|
|
}
|
|
offset = r.offset + offset
|
|
case io.SeekEnd:
|
|
offset = r.block.Size() + offset
|
|
default:
|
|
return 0, errors.New("Seek: invalid whence")
|
|
}
|
|
|
|
if offset < 0 || offset > r.block.Size() {
|
|
return 0, errors.New("Seek: invalid offset")
|
|
}
|
|
r.offset = offset
|
|
return offset, nil
|
|
}
|
|
|
|
func NewDynamicReadAtSeeker(block Block) *DynamicReadAtSeeker {
|
|
return &DynamicReadAtSeeker{
|
|
block: block,
|
|
}
|
|
}
|