internal/packlimit/watch.go

e6cd75b5f28bacf51620bb531320c30fd4e66bfd
gitbay/internal/packlimit/watch.go history · blame · raw

122 lines · 2991 bytes

9 symbols in this file
  1package packlimit
  2
  3import (
  4	"io"
  5	"sync"
  6	"sync/atomic"
  7	"time"
  8)
  9
 10// StallDeadline is how long a limited transport may go without
 11// completing a write to its client before it is killed. upload-pack
 12// sends a keepalive every five seconds while it prepares a pack.
 13var StallDeadline = 2 * time.Minute
 14
 15// Watch wraps w, a transport's writer to its client, so that a client
 16// that stops reading does not hold its slot for as long as its
 17// connection stays open. stalled closes once no write to w has
 18// completed for StallDeadline; stop ends the watch. With no limit in
 19// force (a nil Limiter) nothing is watched: w comes back as is and
 20// stalled never closes.
 21func (l *Limiter) Watch(w io.Writer) (out io.Writer, stalled <-chan struct{}, stop func()) {
 22	if l == nil {
 23		return w, nil, func() {}
 24	}
 25	deadline := StallDeadline
 26	pw := &progressWriter{w: w}
 27	pw.last.Store(time.Now().UnixNano())
 28	st := make(chan struct{})
 29	quit := make(chan struct{})
 30	go func() {
 31		t := time.NewTicker(deadline / 4)
 32		defer t.Stop()
 33		for {
 34			select {
 35			case <-quit:
 36				return
 37			case <-t.C:
 38				if time.Since(time.Unix(0, pw.last.Load())) >= deadline {
 39					close(st)
 40					return
 41				}
 42			}
 43		}
 44	}()
 45	var once sync.Once
 46	return pw, st, func() { once.Do(func() { close(quit) }) }
 47}
 48
 49// progressWriter records when a write to the client last completed.
 50type progressWriter struct {
 51	w    io.Writer
 52	last atomic.Int64 // unix nanoseconds
 53}
 54
 55func (p *progressWriter) Write(b []byte) (int, error) {
 56	n, err := p.w.Write(b)
 57	if n > 0 {
 58		p.last.Store(time.Now().UnixNano())
 59	}
 60	return n, err
 61}
 62
 63// Idle wraps both directions of a transport, r from the client and w to
 64// it, so that once arm has been called idle closes when no byte has
 65// moved either way for d. The window starts at arm; before it idle
 66// never closes. stop ends the watch.
 67func Idle(r io.Reader, w io.Writer, d time.Duration) (in io.Reader, out io.Writer, idle <-chan struct{}, arm, stop func()) {
 68	var last atomic.Int64
 69	var armed atomic.Bool
 70	st := make(chan struct{})
 71	quit := make(chan struct{})
 72	go func() {
 73		t := time.NewTicker(d / 4)
 74		defer t.Stop()
 75		for {
 76			select {
 77			case <-quit:
 78				return
 79			case <-t.C:
 80				if armed.Load() && time.Since(time.Unix(0, last.Load())) >= d {
 81					close(st)
 82					return
 83				}
 84			}
 85		}
 86	}()
 87	var armOnce, stopOnce sync.Once
 88	return &progressReader{r: r, last: &last}, &stampWriter{w: w, last: &last}, st,
 89		func() {
 90			armOnce.Do(func() {
 91				last.Store(time.Now().UnixNano())
 92				armed.Store(true)
 93			})
 94		},
 95		func() { stopOnce.Do(func() { close(quit) }) }
 96}
 97
 98type progressReader struct {
 99	r    io.Reader
100	last *atomic.Int64
101}
102
103func (p *progressReader) Read(b []byte) (int, error) {
104	n, err := p.r.Read(b)
105	if n > 0 {
106		p.last.Store(time.Now().UnixNano())
107	}
108	return n, err
109}
110
111type stampWriter struct {
112	w    io.Writer
113	last *atomic.Int64
114}
115
116func (s *stampWriter) Write(b []byte) (int, error) {
117	n, err := s.w.Write(b)
118	if n > 0 {
119		s.last.Store(time.Now().UnixNano())
120	}
121	return n, err
122}