MCPcopy Create free account
hub / github.com/eapache/channels / batchingBuffer

Method batchingBuffer

batching_channel.go:54–87  ·  view source on GitHub ↗
()

Source from the content-addressed store, hash-verified

52}
53
54func (ch *BatchingChannel) batchingBuffer() {
55 var input, output, nextInput chan interface{}
56 nextInput = ch.input
57 input = nextInput
58
59 for input != nil || output != nil {
60 select {
61 case elem, open := <-input:
62 if open {
63 ch.buffer = append(ch.buffer, elem)
64 } else {
65 input = nil
66 nextInput = nil
67 }
68 case output <- ch.buffer:
69 ch.buffer = nil
70 case ch.length <- len(ch.buffer):
71 }
72
73 if len(ch.buffer) == 0 {
74 input = nextInput
75 output = nil
76 } else if ch.size != Infinity && len(ch.buffer) >= int(ch.size) {
77 input = nil
78 output = ch.output
79 } else {
80 input = nextInput
81 output = ch.output
82 }
83 }
84
85 close(ch.output)
86 close(ch.length)
87}

Callers 1

NewBatchingChannelFunction · 0.95

Calls

no outgoing calls

Tested by

no test coverage detected