internal/packlimit/watch.go
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}