internal/packlimit/watch.go

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

61 lines · 1579 bytes

 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}