cubefs
109 строк · 2.6 Кб
1// Copyright 2023 The CubeFS Authors.
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7// http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or
12// implied. See the License for the specific language governing
13// permissions and limitations under the License.
14
15package unboundedchan
16
17import "sync/atomic"
18
19type UnboundedChan struct {
20bufCount int64
21In chan<- V // for user to write data into, BUT NOT TO READ FROM! In channel can be "close" by user
22Out <-chan V // for user to read data from, BUT NOT TO WRITE INTO!
23buffer *RingBuffer
24}
25
26func (uc *UnboundedChan) Len() int {
27return len(uc.In) + uc.BufLen() + len(uc.Out)
28}
29
30func (uc *UnboundedChan) BufLen() int {
31return int(atomic.LoadInt64(&uc.bufCount))
32}
33
34func NewUnboundedChan(initCapacity int) *UnboundedChan {
35return NewUnboundedChanSize(initCapacity, initCapacity, initCapacity)
36}
37
38func NewUnboundedChanSize(initInCapacity, initOutCapacity, initBufferCapacity int) *UnboundedChan {
39in := make(chan V, initInCapacity)
40out := make(chan V, initOutCapacity)
41uc := &UnboundedChan{
42In: in,
43Out: out,
44buffer: NewRingBuffer(uint32(initBufferCapacity)),
45}
46
47go run(in, out, uc)
48return uc
49}
50
51func (uc *UnboundedChan) feedBuffer(val V) {
52uc.buffer.Write(val)
53atomic.AddInt64(&uc.bufCount, 1)
54}
55
56func (uc *UnboundedChan) resetBuffer() {
57uc.buffer.Reset()
58atomic.StoreInt64(&uc.bufCount, 0)
59}
60
61func run(in, out chan V, uc *UnboundedChan) {
62defer close(out)
63LOOP:
64for {
65// read data from in channel
66val, ok := <-in
67if !ok {
68break LOOP
69}
70
71// move data to buffer or out channel
72if atomic.LoadInt64(&uc.bufCount) > 0 {
73uc.feedBuffer(val)
74} else {
75select {
76case out <- val:
77continue
78default:
79// out channel is full
80uc.feedBuffer(val)
81}
82}
83
84// try feeding out channel with buffer data
85for !uc.buffer.IsEmpty() {
86select {
87// when out channel is full, data may still feed in channel
88case val, ok := <-in:
89if !ok {
90break LOOP
91}
92uc.feedBuffer(val)
93case out <- uc.buffer.Peek():
94uc.buffer.Pop()
95atomic.AddInt64(&uc.bufCount, -1)
96if uc.buffer.IsEmpty() {
97uc.resetBuffer()
98}
99}
100}
101}
102
103// no more in data, keep feeding out channel with buffer data until it's drained
104for !uc.buffer.IsEmpty() {
105out <- uc.buffer.Pop()
106atomic.AddInt64(&uc.bufCount, -1)
107}
108uc.resetBuffer()
109}
110