stream

package
v4.3.5 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Aug 7, 2025 License: AGPL-3.0 Imports: 20 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

View Source
var (
	ErrInvalidLink = errs.New("无效的链接")
)

错误定义

Functions

func CacheFullInTempFileAndHash

func CacheFullInTempFileAndHash(stream model.FileStreamer, up model.UpdateProgress, hashType *utils.HashType, hashParams ...any) (model.File, string, error)

CacheFullInTempFileAndHash 将流的内容缓存到临时文件并计算哈希值 如果提供了进度更新函数,则在读取过程中更新进度

func CacheFullInTempFileAndWriter

func CacheFullInTempFileAndWriter(stream model.FileStreamer, up model.UpdateProgress, w io.Writer) (model.File, error)

CacheFullInTempFileAndWriter 将流的内容缓存到临时文件并同时写入指定的写入器 如果流已经有缓存文件,则直接使用该文件 如果提供了进度更新函数,则在读取过程中更新进度

func GetRangeReaderFromLink(size int64, link *model.Link) (model.RangeReaderIF, error)

GetRangeReaderFromLink 从链接创建范围读取器 它支持多种类型的链接,包括文件、URL和自定义范围读取器

func GetRangeReaderFromMFile added in v4.2.3

func GetRangeReaderFromMFile(size int64, file model.File) model.RangeReaderIF

GetRangeReaderFromMFile RangeReaderIF.RangeRead返回的io.ReadCloser保留file的签名。

func NewMultiReaderAt

func NewMultiReaderAt(ss []*SeekableStream) (readerutil.SizeReaderAt, error)

func NewReadAtSeeker

func NewReadAtSeeker(ss *SeekableStream, offset int64, forceRange ...bool) (model.File, error)

Types

type FileStream

type FileStream struct {
	Ctx context.Context
	model.Obj
	io.Reader
	Mimetype          string
	WebPutAsTask      bool
	ForceStreamUpload bool
	Exist             model.Obj // the file existed in the destination, we can reuse some info since we wil overwrite it
	utils.Closers
	// contains filtered or unexported fields
}

func (*FileStream) CacheFullInTempFile

func (f *FileStream) CacheFullInTempFile() (model.File, error)

CacheFullInTempFile save all data into tmpFile. Not recommended since it wears disk, and can't start upload until the file is written. It's not thread-safe!

func (*FileStream) Close

func (f *FileStream) Close() error

func (*FileStream) GetExist

func (f *FileStream) GetExist() model.Obj

func (*FileStream) GetFile

func (f *FileStream) GetFile() model.File

func (*FileStream) GetMimetype

func (f *FileStream) GetMimetype() string

func (*FileStream) GetSize

func (f *FileStream) GetSize() int64

func (*FileStream) IsForceStreamUpload

func (f *FileStream) IsForceStreamUpload() bool

func (*FileStream) NeedStore

func (f *FileStream) NeedStore() bool

func (*FileStream) RangeRead

func (f *FileStream) RangeRead(httpRange http_range.Range) (io.Reader, error)

RangeRead have to cache all data first since only Reader is provided. also support a peeking RangeRead at very start, but won't buffer more than conf.MaxBufferLimit data in memory

func (*FileStream) SetExist

func (f *FileStream) SetExist(obj model.Obj)

func (*FileStream) SetTmpFile

func (f *FileStream) SetTmpFile(r *os.File)

type Limiter

type Limiter interface {
	// 基本速率限制方法
	Limit() rate.Limit
	Burst() int
	TokensAt(time.Time) float64
	Tokens() float64
	Allow() bool
	AllowN(time.Time, int) bool
	Reserve() *rate.Reservation
	ReserveN(time.Time, int) *rate.Reservation
	Wait(context.Context) error
	WaitN(context.Context, int) error

	// 设置限制方法
	SetLimit(rate.Limit)
	SetLimitAt(time.Time, rate.Limit)
	SetBurst(int)
	SetBurstAt(time.Time, int)
}

Limiter 接口定义了速率限制器的行为 它扩展了 golang.org/x/time/rate.Limiter 接口,添加了一些额外的方法

var (
	// ClientDownloadLimit 客户端下载速率限制器
	ClientDownloadLimit Limiter
	// ClientUploadLimit 客户端上传速率限制器
	ClientUploadLimit Limiter
	// ServerDownloadLimit 服务器下载速率限制器
	ServerDownloadLimit Limiter
	// ServerUploadLimit 服务器上传速率限制器
	ServerUploadLimit Limiter
)

全局速率限制器

type RangeReadReadAtSeeker

type RangeReadReadAtSeeker struct {
	// contains filtered or unexported fields
}

func (*RangeReadReadAtSeeker) InitHeadCache

func (r *RangeReadReadAtSeeker) InitHeadCache()

func (*RangeReadReadAtSeeker) Read

func (r *RangeReadReadAtSeeker) Read(p []byte) (n int, err error)

func (*RangeReadReadAtSeeker) ReadAt

func (r *RangeReadReadAtSeeker) ReadAt(p []byte, off int64) (n int, err error)

func (*RangeReadReadAtSeeker) Seek

func (r *RangeReadReadAtSeeker) Seek(offset int64, whence int) (int64, error)

type RangeReaderFunc added in v4.2.3

type RangeReaderFunc func(ctx context.Context, httpRange http_range.Range) (io.ReadCloser, error)

RangeReaderFunc 是一个函数类型,用于实现 model.RangeReaderIF 接口 它允许将普通函数转换为范围读取器

func (RangeReaderFunc) RangeRead added in v4.2.3

func (f RangeReaderFunc) RangeRead(ctx context.Context, httpRange http_range.Range) (io.ReadCloser, error)

RangeRead 实现 model.RangeReaderIF 接口 它调用底层函数来执行范围读取

type RateLimitFile

type RateLimitFile struct {
	model.File                 // 底层文件
	Limiter    Limiter         // 速率限制器
	Ctx        context.Context // 上下文,用于取消操作
}

RateLimitFile 实现了一个带速率限制的文件接口 它在每次读取操作后等待适当的时间,以确保不超过指定的速率

func (*RateLimitFile) Close added in v4.1.8

func (r *RateLimitFile) Close() error

Close 实现了 io.Closer 接口

func (*RateLimitFile) Read

func (r *RateLimitFile) Read(p []byte) (n int, err error)

Read 实现了 io.Reader 接口,增加了速率限制

func (*RateLimitFile) ReadAt

func (r *RateLimitFile) ReadAt(p []byte, off int64) (n int, err error)

ReadAt 实现了 io.ReaderAt 接口,增加了速率限制

type RateLimitRangeReaderFunc added in v4.2.3

type RateLimitRangeReaderFunc RangeReaderFunc

RateLimitRangeReaderFunc 是一个带速率限制的范围读取函数

func (RateLimitRangeReaderFunc) RangeRead added in v4.2.3

func (f RateLimitRangeReaderFunc) RangeRead(ctx context.Context, httpRange http_range.Range) (io.ReadCloser, error)

RangeRead 实现了 model.RangeReaderIF 接口,增加了速率限制 它首先调用底层的范围读取函数,然后将结果包装在一个带速率限制的读取器中

type RateLimitReader

type RateLimitReader struct {
	io.Reader                 // 底层读取器
	Limiter   Limiter         // 速率限制器
	Ctx       context.Context // 上下文,用于取消操作
}

RateLimitReader 实现了一个带速率限制的读取器 它在每次读取操作后等待适当的时间,以确保不超过指定的速率

func (*RateLimitReader) Close

func (r *RateLimitReader) Close() error

Close 实现了 io.Closer 接口 如果底层读取器支持关闭,则关闭它

func (*RateLimitReader) Read

func (r *RateLimitReader) Read(p []byte) (n int, err error)

Read 实现了 io.Reader 接口,增加了速率限制 它首先检查上下文是否已取消,然后从底层读取器读取数据, 最后等待足够的时间以确保不超过速率限制

type RateLimitWriter

type RateLimitWriter struct {
	io.Writer                 // 底层写入器
	Limiter   Limiter         // 速率限制器
	Ctx       context.Context // 上下文,用于取消操作
}

RateLimitWriter 实现了一个带速率限制的写入器 它在每次写入操作后等待适当的时间,以确保不超过指定的速率

func (*RateLimitWriter) Close

func (w *RateLimitWriter) Close() error

Close 实现了 io.Closer 接口 如果底层写入器支持关闭,则关闭它

func (*RateLimitWriter) Write

func (w *RateLimitWriter) Write(p []byte) (n int, err error)

Write 实现了 io.Writer 接口,增加了速率限制 它首先检查上下文是否已取消,然后向底层写入器写入数据, 最后等待足够的时间以确保不超过速率限制

type ReaderUpdatingProgress

type ReaderUpdatingProgress struct {
	Reader ReaderWithSize
	model.UpdateProgress
	// contains filtered or unexported fields
}

func (*ReaderUpdatingProgress) Close

func (r *ReaderUpdatingProgress) Close() error

func (*ReaderUpdatingProgress) Read

func (r *ReaderUpdatingProgress) Read(p []byte) (n int, err error)

type ReaderWithCtx

type ReaderWithCtx struct {
	io.Reader
	Ctx context.Context
}

ReaderWithCtx 是一个带有上下文的读取器 它在每次读取操作前检查上下文是否已取消

func (*ReaderWithCtx) Close

func (r *ReaderWithCtx) Close() error

Close 实现io.Closer接口

func (*ReaderWithCtx) Read

func (r *ReaderWithCtx) Read(p []byte) (n int, err error)

Read 实现io.Reader接口,增加了上下文取消检查

type ReaderWithSize

type ReaderWithSize interface {
	io.ReadCloser
	GetSize() int64
}

type SectionReader added in v4.3.4

type SectionReader struct {
	io.ReadSeeker
	// contains filtered or unexported fields
}

type SeekableStream

type SeekableStream struct {
	*FileStream
	// contains filtered or unexported fields
}

SeekableStream for most internal stream, which is either RangeReadCloser or MFile Any functionality implemented based on SeekableStream should implement a Close method, whose only purpose is to close the SeekableStream object. If such functionality has additional resources that need to be closed, they should be added to the Closer property of the SeekableStream object and be closed together when the SeekableStream object is closed.

func NewSeekableStream

func NewSeekableStream(fs *FileStream, link *model.Link) (*SeekableStream, error)

func (*SeekableStream) CacheFullInTempFile

func (ss *SeekableStream) CacheFullInTempFile() (model.File, error)

func (*SeekableStream) GetSize added in v4.2.4

func (ss *SeekableStream) GetSize() int64

func (*SeekableStream) RangeRead

func (ss *SeekableStream) RangeRead(httpRange http_range.Range) (io.Reader, error)

RangeRead is not thread-safe, pls use it in single thread only.

func (*SeekableStream) Read

func (ss *SeekableStream) Read(p []byte) (n int, err error)

only provide Reader as full stream when it's demanded. in rapid-upload, we can skip this to save memory

type SimpleReaderWithSize

type SimpleReaderWithSize struct {
	io.Reader
	Size int64
}

func (*SimpleReaderWithSize) Close

func (r *SimpleReaderWithSize) Close() error

func (*SimpleReaderWithSize) GetSize

func (r *SimpleReaderWithSize) GetSize() int64

type StreamSectionReader added in v4.3.4

type StreamSectionReader struct {
	// contains filtered or unexported fields
}

func NewStreamSectionReader added in v4.3.4

func NewStreamSectionReader(file model.FileStreamer, maxBufferSize int) (*StreamSectionReader, error)

func (*StreamSectionReader) GetSectionReader added in v4.3.4

func (ss *StreamSectionReader) GetSectionReader(off, length int64) (*SectionReader, error)

func (*StreamSectionReader) RecycleSectionReader added in v4.3.4

func (ss *StreamSectionReader) RecycleSectionReader(sr *SectionReader)

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL