package stario import ( "errors" "fmt" "io" "os" "runtime" "sync" "sync/atomic" "time" ) type StarBuffer struct { io.Reader io.Writer io.Closer datas []byte pStart uint64 pEnd uint64 cap uint64 isClose atomic.Value isEnd atomic.Value rmu sync.Mutex wmu sync.Mutex } func NewStarBuffer(cap uint64) *StarBuffer { rtnBuffer := new(StarBuffer) rtnBuffer.cap = cap rtnBuffer.datas = make([]byte, cap) rtnBuffer.isClose.Store(false) rtnBuffer.isEnd.Store(false) return rtnBuffer } func (star *StarBuffer) Free() uint64 { return star.cap - star.Len() } func (star *StarBuffer) Cap() uint64 { return star.cap } func (star *StarBuffer) Len() uint64 { if star.pEnd >= star.pStart { return star.pEnd - star.pStart } return star.pEnd - star.pStart + star.cap } func (star *StarBuffer) getByte() (byte, error) { if star.isClose.Load().(bool) || (star.Len() == 0 && star.isEnd.Load().(bool)) { return 0, io.EOF } if star.Len() == 0 { return 0, os.ErrNotExist } nowPtr := star.pStart nextPtr := star.pStart + 1 if nextPtr >= star.cap { nextPtr = 0 } data := star.datas[nowPtr] ok := atomic.CompareAndSwapUint64(&star.pStart, nowPtr, nextPtr) if !ok { return 0, os.ErrInvalid } return data, nil } func (star *StarBuffer) putByte(data byte) error { if star.isClose.Load().(bool) { return io.EOF } nowPtr := star.pEnd kariEnd := nowPtr + 1 if kariEnd == star.cap { kariEnd = 0 } if kariEnd == atomic.LoadUint64(&star.pStart) { for { time.Sleep(time.Microsecond) runtime.Gosched() if kariEnd != atomic.LoadUint64(&star.pStart) { break } } } star.datas[nowPtr] = data if ok := atomic.CompareAndSwapUint64(&star.pEnd, nowPtr, kariEnd); !ok { return os.ErrInvalid } return nil } func (star *StarBuffer) Close() error { star.isClose.Store(true) return nil } func (star *StarBuffer) Read(buf []byte) (int, error) { if star.isClose.Load().(bool) || (star.Len() == 0 && star.isEnd.Load().(bool)) { return 0, io.EOF } if buf == nil { return 0, errors.New("buffer is nil") } star.rmu.Lock() defer star.rmu.Unlock() var sum int = 0 for i := 0; i < len(buf); i++ { data, err := star.getByte() if err != nil { if err == io.EOF { return sum, err } if err == os.ErrNotExist { i-- continue } return sum, nil } buf[i] = data sum++ } return sum, nil } func (star *StarBuffer) Write(bts []byte) (int, error) { if bts == nil && !star.isEnd.Load().(bool) { star.isEnd.Store(true) return 0, nil } if bts == nil || star.isClose.Load().(bool) { return 0, io.EOF } star.wmu.Lock() defer star.wmu.Unlock() var sum = 0 for i := 0; i < len(bts); i++ { err := star.putByte(bts[i]) if err != nil { fmt.Println("Write bts err:", err) return sum, err } sum++ } return sum, nil }