Commit 03e34fb24d
Verified · cmc
Layout: unified · split
internal/config/config.go +28
| @@ -43,6 +43,11 @@ const ( | |||
| 43 | DefaultPushPerPrincipal = 1 | 43 | DefaultPushPerPrincipal = 1 |
| 44 | DefaultPushQueue = 16 | 44 | DefaultPushQueue = 16 |
| 45 | DefaultPushQueueWait = time.Minute | 45 | DefaultPushQueueWait = time.Minute |
| 46 | // A push is killed when nothing has moved either way for | ||
| 47 | // DefaultPushIdle, or when its pre-receive has not started | ||
| 48 | // DefaultPushReceiveTimeout after it took its slot. | ||
| 49 | DefaultPushIdle = time.Minute | ||
| 50 | DefaultPushReceiveTimeout = 15 * time.Minute | ||
| 46 | ) | 51 | ) |
| 47 | 52 | ||
| 48 | type Config struct { | 53 | type Config struct { |
| @@ -253,6 +258,15 @@ type Limits struct { | |||
| 253 | PushPerPrincipal int `toml:"push_per_principal"` | 258 | PushPerPrincipal int `toml:"push_per_principal"` |
| 254 | PushQueue int `toml:"push_queue"` | 259 | PushQueue int `toml:"push_queue"` |
| 255 | PushQueueWait string `toml:"push_queue_wait"` | 260 | PushQueueWait string `toml:"push_queue_wait"` |
| 261 | // PushIdle ("60s") kills a push when no byte has been read from the | ||
| 262 | // client and none written to it for that long, counted from when | ||
| 263 | // the pack begins or pre-receive starts; receive-pack sends a | ||
| 264 | // keepalive while it indexes and runs hooks. PushReceiveTimeout | ||
| 265 | // ("15m") kills a push whose pre-receive has not started that long | ||
| 266 | // after it took its slot, bounding how long a pack may take to | ||
| 267 | // arrive. Both apply only where the push limit is in force. | ||
| 268 | PushIdle string `toml:"push_idle"` | ||
| 269 | PushReceiveTimeout string `toml:"push_receive_timeout"` | ||
| 256 | } | 270 | } |
| 257 | 271 | ||
| 258 | // PackLimits resolves the pack_* settings for packlimit.New. A zero | 272 | // PackLimits resolves the pack_* settings for packlimit.New. A zero |
| @@ -270,6 +284,18 @@ func (l Limits) PushLimits() (max, per, queue int, wait time.Duration) { | |||
| 270 | DefaultPushConcurrency, DefaultPushPerPrincipal, DefaultPushQueue, DefaultPushQueueWait) | 284 | DefaultPushConcurrency, DefaultPushPerPrincipal, DefaultPushQueue, DefaultPushQueueWait) |
| 271 | } | 285 | } |
| 272 | 286 | ||
| 287 | // PushTimeouts resolves push_idle and push_receive_timeout. | ||
| 288 | func (l Limits) PushTimeouts() (idle, receive time.Duration) { | ||
| 289 | idle, receive = DefaultPushIdle, DefaultPushReceiveTimeout | ||
| 290 | if d, err := time.ParseDuration(l.PushIdle); err == nil && d > 0 { | ||
| 291 | idle = d | ||
| 292 | } | ||
| 293 | if d, err := time.ParseDuration(l.PushReceiveTimeout); err == nil && d > 0 { | ||
| 294 | receive = d | ||
| 295 | } | ||
| 296 | return idle, receive | ||
| 297 | } | ||
| 298 | |||
| 273 | func resolveLimits(maxV, perV, queueV int, waitV string, maxDef, perDef, queueDef int, waitDef time.Duration) (max, per, queue int, wait time.Duration) { | 299 | func resolveLimits(maxV, perV, queueV int, waitV string, maxDef, perDef, queueDef int, waitDef time.Duration) (max, per, queue int, wait time.Duration) { |
| 274 | pick := func(v, def int) int { | 300 | pick := func(v, def int) int { |
| 275 | switch { | 301 | switch { |
| @@ -651,6 +677,8 @@ func (c Config) Validate() error { | |||
| 651 | for _, w := range []struct{ name, val string }{ | 677 | for _, w := range []struct{ name, val string }{ |
| 652 | {"pack_queue_wait", c.Limits.PackQueueWait}, | 678 | {"pack_queue_wait", c.Limits.PackQueueWait}, |
| 653 | {"push_queue_wait", c.Limits.PushQueueWait}, | 679 | {"push_queue_wait", c.Limits.PushQueueWait}, |
| 680 | {"push_idle", c.Limits.PushIdle}, | ||
| 681 | {"push_receive_timeout", c.Limits.PushReceiveTimeout}, | ||
| 654 | } { | 682 | } { |
| 655 | if w.val == "" { | 683 | if w.val == "" { |
| 656 | continue | 684 | continue |
internal/config/config_test.go +17 −1
| @@ -76,7 +76,7 @@ func TestPushLimits(t *testing.T) { | |||
| 76 | if max != 0 || per != 0 || queue != math.MaxInt || wait != 5*time.Second { | 76 | if max != 0 || per != 0 || queue != math.MaxInt || wait != 5*time.Second { |
| 77 | t.Fatalf("off: %d %d %d %s", max, per, queue, wait) | 77 | t.Fatalf("off: %d %d %d %s", max, per, queue, wait) |
| 78 | } | 78 | } |
| 79 | cfg, err := Load(writeConfig(t, minimal+"\n[limits]\npush_concurrency = 4\npush_per_principal = 2\npush_queue = 8\npush_queue_wait = \"30s\"\n")) | 79 | cfg, err := Load(writeConfig(t, minimal+"\n[limits]\npush_concurrency = 4\npush_per_principal = 2\npush_queue = 8\npush_queue_wait = \"30s\"\npush_idle = \"20s\"\npush_receive_timeout = \"5m\"\n")) |
| 80 | if err != nil { | 80 | if err != nil { |
| 81 | t.Fatal(err) | 81 | t.Fatal(err) |
| 82 | } | 82 | } |
| @@ -84,6 +84,12 @@ func TestPushLimits(t *testing.T) { | |||
| 84 | if max != 4 || per != 2 || queue != 8 || wait != 30*time.Second { | 84 | if max != 4 || per != 2 || queue != 8 || wait != 30*time.Second { |
| 85 | t.Fatalf("loaded: %d %d %d %s", max, per, queue, wait) | 85 | t.Fatalf("loaded: %d %d %d %s", max, per, queue, wait) |
| 86 | } | 86 | } |
| 87 | if idle, receive := cfg.Limits.PushTimeouts(); idle != 20*time.Second || receive != 5*time.Minute { | ||
| 88 | t.Fatalf("timeouts: %s %s", idle, receive) | ||
| 89 | } | ||
| 90 | if idle, receive := (Limits{}).PushTimeouts(); idle != time.Minute || receive != 15*time.Minute { | ||
| 91 | t.Fatalf("default timeouts: %s %s", idle, receive) | ||
| 92 | } | ||
| 87 | // The pack budget is not read from the push settings. | 93 | // The pack budget is not read from the push settings. |
| 88 | if pm, _, _, _ := cfg.Limits.PackLimits(); pm != DefaultPackConcurrency { | 94 | if pm, _, _, _ := cfg.Limits.PackLimits(); pm != DefaultPackConcurrency { |
| 89 | t.Fatalf("pack_concurrency %d, want the default", pm) | 95 | t.Fatalf("pack_concurrency %d, want the default", pm) |
| @@ -106,6 +112,16 @@ func TestContradictions(t *testing.T) { | |||
| 106 | minimal + "\n[limits]\npush_queue_wait = \"-5s\"\n", | 112 | minimal + "\n[limits]\npush_queue_wait = \"-5s\"\n", |
| 107 | "limits.push_queue_wait", | 113 | "limits.push_queue_wait", |
| 108 | }, | 114 | }, |
| 115 | { | ||
| 116 | "bad push_idle", | ||
| 117 | minimal + "\n[limits]\npush_idle = \"0s\"\n", | ||
| 118 | "limits.push_idle", | ||
| 119 | }, | ||
| 120 | { | ||
| 121 | "bad push_receive_timeout", | ||
| 122 | minimal + "\n[limits]\npush_receive_timeout = \"later\"\n", | ||
| 123 | "limits.push_receive_timeout", | ||
| 124 | }, | ||
| 109 | { | 125 | { |
| 110 | "registration open without smtp", | 126 | "registration open without smtp", |
| 111 | minimal + "\n[registration]\nmode = \"open\"\n", | 127 | minimal + "\n[registration]\nmode = \"open\"\n", |
internal/gitd/gitd.go +1 −1
| @@ -99,7 +99,7 @@ func (s *Server) handle(conn net.Conn) { | |||
| 99 | }() | 99 | }() |
| 100 | 100 | ||
| 101 | dir := control.RepoDir(s.cfg.Server.Root, repo.OwnerName, repo.Name) | 101 | dir := control.RepoDir(s.cfg.Server.Root, repo.OwnerName, repo.Name) |
| 102 | gitutil.Transport("git-upload-pack", dir, conn, out, io.Discard, protoEnv, 0, kill) | 102 | gitutil.Transport("git-upload-pack", dir, conn, out, io.Discard, protoEnv, 0, 0, kill) |
| 103 | } | 103 | } |
| 104 | 104 | ||
| 105 | // principal is the pack-limit principal for a client at addr. | 105 | // principal is the pack-limit principal for a client at addr. |
internal/gitutil/gitutil.go +8 −2
| @@ -35,16 +35,22 @@ func InitBare(path, defaultBranch, hooksPath string) error { | |||
| 35 | // Transport streams one git transport service (upload-pack, receive-pack, | 35 | // Transport streams one git transport service (upload-pack, receive-pack, |
| 36 | // upload-archive). extraEnv entries are appended to the process environment; | 36 | // upload-archive). extraEnv entries are appended to the process environment; |
| 37 | // hooks read the GITBAY_* variables from it. maxPack caps incoming pack | 37 | // hooks read the GITBAY_* variables from it. maxPack caps incoming pack |
| 38 | // bytes on receive-pack (0 = unlimited). Closing cancel kills the service | 38 | // bytes on receive-pack (0 = unlimited), and keepAlive sets its |
| 39 | // receive.keepAlive in seconds (0 = git's default). Closing cancel kills the service | ||
| 39 | // and everything it started; a push killed before its pre-receive hook | 40 | // and everything it started; a push killed before its pre-receive hook |
| 40 | // answers updates no refs. A nil cancel never fires. | 41 | // answers updates no refs. A nil cancel never fires. |
| 41 | func Transport(service, repoPath string, stdin io.Reader, stdout, errW io.Writer, extraEnv []string, maxPack int64, cancel <-chan struct{}) error { | 42 | func Transport(service, repoPath string, stdin io.Reader, stdout, errW io.Writer, extraEnv []string, maxPack int64, keepAlive int, cancel <-chan struct{}) error { |
| 42 | var args []string | 43 | var args []string |
| 43 | switch service { | 44 | switch service { |
| 44 | case "git-upload-pack", "git-receive-pack", "git-upload-archive": | 45 | case "git-upload-pack", "git-receive-pack", "git-upload-archive": |
| 45 | if service == "git-receive-pack" && maxPack > 0 { | 46 | if service == "git-receive-pack" && maxPack > 0 { |
| 46 | args = []string{"-c", fmt.Sprintf("receive.maxInputSize=%d", maxPack)} | 47 | args = []string{"-c", fmt.Sprintf("receive.maxInputSize=%d", maxPack)} |
| 47 | } | 48 | } |
| 49 | if service == "git-receive-pack" && keepAlive > 0 { | ||
| 50 | // Seconds of silence after which receive-pack sends a | ||
| 51 | // keepalive while it indexes the pack and runs hooks. | ||
| 52 | args = append(args, "-c", fmt.Sprintf("receive.keepAlive=%d", keepAlive)) | ||
| 53 | } | ||
| 48 | if service == "git-upload-pack" { | 54 | if service == "git-upload-pack" { |
| 49 | // Keepalives while pack-objects is still counting keep a | 55 | // Keepalives while pack-objects is still counting keep a |
| 50 | // healthy clone writing; a limited transport kills one that goes quiet. | 56 | // healthy clone writing; a limited transport kills one that goes quiet. |
internal/gitutil/transport_test.go +1 −1
| @@ -19,7 +19,7 @@ func TestTransportCancelKillsGit(t *testing.T) { | |||
| 19 | in, w := io.Pipe() | 19 | in, w := io.Pipe() |
| 20 | cancel := make(chan struct{}) | 20 | cancel := make(chan struct{}) |
| 21 | errc := make(chan error, 1) | 21 | errc := make(chan error, 1) |
| 22 | go func() { errc <- Transport("git-upload-pack", dir, in, io.Discard, io.Discard, nil, 0, cancel) }() | 22 | go func() { errc <- Transport("git-upload-pack", dir, in, io.Discard, io.Discard, nil, 0, 0, cancel) }() |
| 23 | close(cancel) | 23 | close(cancel) |
| 24 | time.AfterFunc(500*time.Millisecond, func() { w.Close() }) | 24 | time.AfterFunc(500*time.Millisecond, func() { w.Close() }) |
| 25 | select { | 25 | select { |
internal/hookd/hookd.go +35
| @@ -17,6 +17,7 @@ import ( | |||
| 17 | "os" | 17 | "os" |
| 18 | "path/filepath" | 18 | "path/filepath" |
| 19 | "strings" | 19 | "strings" |
| 20 | "sync" | ||
| 20 | 21 | ||
| 21 | "gitbay.org/gitbay/internal/config" | 22 | "gitbay.org/gitbay/internal/config" |
| 22 | "gitbay.org/gitbay/internal/control" | 23 | "gitbay.org/gitbay/internal/control" |
| @@ -139,6 +140,7 @@ func (s *Server) handle(conn net.Conn) { | |||
| 139 | } | 140 | } |
| 140 | switch req.Hook { | 141 | switch req.Hook { |
| 141 | case "pre-receive": | 142 | case "pre-receive": |
| 143 | preReceiveStarted(req.Token) | ||
| 142 | s.preReceive(req, dec, enc) | 144 | s.preReceive(req, dec, enc) |
| 143 | case "post-receive": | 145 | case "post-receive": |
| 144 | s.postReceive(req) | 146 | s.postReceive(req) |
| @@ -324,3 +326,36 @@ func WriteHookScripts(hooksDir, gitbaydPath string) error { | |||
| 324 | } | 326 | } |
| 325 | return nil | 327 | return nil |
| 326 | } | 328 | } |
| 329 | |||
| 330 | // awaiting holds, per push token, a channel closed when that push's | ||
| 331 | // pre-receive reaches hookd. sshd uses it to tell a pack still arriving | ||
| 332 | // from one being checked. | ||
| 333 | var ( | ||
| 334 | awaitMu sync.Mutex | ||
| 335 | awaiting = map[string]chan struct{}{} | ||
| 336 | ) | ||
| 337 | |||
| 338 | // AwaitPreReceive returns a channel that closes when pre-receive for | ||
| 339 | // the push holding token reaches this process's hookd, and forget, | ||
| 340 | // which drops the registration. Under ssh.mode = "system" the push runs | ||
| 341 | // in another process and the channel never closes. | ||
| 342 | func AwaitPreReceive(token string) (started <-chan struct{}, forget func()) { | ||
| 343 | ch := make(chan struct{}) | ||
| 344 | awaitMu.Lock() | ||
| 345 | awaiting[token] = ch | ||
| 346 | awaitMu.Unlock() | ||
| 347 | return ch, func() { | ||
| 348 | awaitMu.Lock() | ||
| 349 | delete(awaiting, token) | ||
| 350 | awaitMu.Unlock() | ||
| 351 | } | ||
| 352 | } | ||
| 353 | |||
| 354 | func preReceiveStarted(token string) { | ||
| 355 | awaitMu.Lock() | ||
| 356 | defer awaitMu.Unlock() | ||
| 357 | if ch, ok := awaiting[token]; ok { | ||
| 358 | close(ch) | ||
| 359 | delete(awaiting, token) | ||
| 360 | } | ||
| 361 | } | ||
internal/hookd/socket_test.go +31
| @@ -9,6 +9,7 @@ import ( | |||
| 9 | "path/filepath" | 9 | "path/filepath" |
| 10 | "strings" | 10 | "strings" |
| 11 | "testing" | 11 | "testing" |
| 12 | "time" | ||
| 12 | 13 | ||
| 13 | "gitbay.org/gitbay/internal/config" | 14 | "gitbay.org/gitbay/internal/config" |
| 14 | "gitbay.org/gitbay/internal/policy" | 15 | "gitbay.org/gitbay/internal/policy" |
| @@ -198,3 +199,33 @@ func TestRefusedPushCapsRefs(t *testing.T) { | |||
| 198 | t.Fatalf("refs %d, more %d", len(data.Refs), data.MoreRefs) | 199 | t.Fatalf("refs %d, more %d", len(data.Refs), data.MoreRefs) |
| 199 | } | 200 | } |
| 200 | } | 201 | } |
| 202 | |||
| 203 | // A push's pre-receive reaching hookd closes the channel sshd waits on; | ||
| 204 | // a request that fails the token check does not. | ||
| 205 | func TestPreReceiveSignalsItsPush(t *testing.T) { | ||
| 206 | sock, st, repoID, uid := serveSocket(t) | ||
| 207 | token, err := st.CreatePushToken(repoID, uid, "full") | ||
| 208 | if err != nil { | ||
| 209 | t.Fatal(err) | ||
| 210 | } | ||
| 211 | started, forget := AwaitPreReceive(token) | ||
| 212 | defer forget() | ||
| 213 | forged := Request{Hook: "pre-receive", RepoID: repoID, UserID: uid, Scope: "read", Token: token} | ||
| 214 | if resp, err := Ask(sock, forged, nil); err != nil || resp.Allow { | ||
| 215 | t.Fatalf("forged: %+v, %v", resp, err) | ||
| 216 | } | ||
| 217 | select { | ||
| 218 | case <-started: | ||
| 219 | t.Fatal("a refused request signalled the push") | ||
| 220 | default: | ||
| 221 | } | ||
| 222 | req := Request{Hook: "pre-receive", RepoID: repoID, UserID: uid, Scope: "full", Token: token} | ||
| 223 | if resp, err := Ask(sock, req, nil); err != nil || !resp.Allow { | ||
| 224 | t.Fatalf("pre-receive: %+v, %v", resp, err) | ||
| 225 | } | ||
| 226 | select { | ||
| 227 | case <-started: | ||
| 228 | case <-time.After(5 * time.Second): | ||
| 229 | t.Fatal("pre-receive did not signal the push") | ||
| 230 | } | ||
| 231 | } | ||
internal/packlimit/watch.go +61
| @@ -59,3 +59,64 @@ func (p *progressWriter) Write(b []byte) (int, error) { | |||
| 59 | } | 59 | } |
| 60 | return n, err | 60 | return n, err |
| 61 | } | 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. | ||
| 67 | func 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 | |||
| 98 | type progressReader struct { | ||
| 99 | r io.Reader | ||
| 100 | last *atomic.Int64 | ||
| 101 | } | ||
| 102 | |||
| 103 | func (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 | |||
| 111 | type stampWriter struct { | ||
| 112 | w io.Writer | ||
| 113 | last *atomic.Int64 | ||
| 114 | } | ||
| 115 | |||
| 116 | func (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 | } | ||
internal/packlimit/watch_test.go +35
| @@ -2,6 +2,7 @@ package packlimit | |||
| 2 | 2 | ||
| 3 | import ( | 3 | import ( |
| 4 | "io" | 4 | "io" |
| 5 | "strings" | ||
| 5 | "testing" | 6 | "testing" |
| 6 | "time" | 7 | "time" |
| 7 | ) | 8 | ) |
| @@ -38,3 +39,37 @@ func TestWatchWithoutLimit(t *testing.T) { | |||
| 38 | t.Fatal("a nil limiter watched") | 39 | t.Fatal("a nil limiter watched") |
| 39 | } | 40 | } |
| 40 | } | 41 | } |
| 42 | |||
| 43 | // Bytes either way keep an armed Idle watch alive; silence both ways | ||
| 44 | // ends it. Before arm, silence ends nothing. | ||
| 45 | func TestIdleCountsBothDirections(t *testing.T) { | ||
| 46 | src := strings.NewReader(strings.Repeat("x", 16)) | ||
| 47 | r, w, idle, arm, stop := Idle(src, io.Discard, 100*time.Millisecond) | ||
| 48 | defer stop() | ||
| 49 | time.Sleep(300 * time.Millisecond) | ||
| 50 | select { | ||
| 51 | case <-idle: | ||
| 52 | t.Fatal("idle before arm") | ||
| 53 | default: | ||
| 54 | } | ||
| 55 | arm() | ||
| 56 | buf := make([]byte, 1) | ||
| 57 | for i := range 8 { | ||
| 58 | time.Sleep(40 * time.Millisecond) | ||
| 59 | if i%2 == 0 { | ||
| 60 | r.Read(buf) | ||
| 61 | } else { | ||
| 62 | w.Write(buf) | ||
| 63 | } | ||
| 64 | select { | ||
| 65 | case <-idle: | ||
| 66 | t.Fatalf("idle at step %d while bytes moved", i) | ||
| 67 | default: | ||
| 68 | } | ||
| 69 | } | ||
| 70 | select { | ||
| 71 | case <-idle: | ||
| 72 | case <-time.After(2 * time.Second): | ||
| 73 | t.Fatal("not idle after both directions went quiet") | ||
| 74 | } | ||
| 75 | } | ||
internal/sshd/pushlimit_test.go added +275
| @@ -0,0 +1,275 @@ | |||
| 1 | package sshd | ||
| 2 | |||
| 3 | import ( | ||
| 4 | "crypto/ed25519" | ||
| 5 | "crypto/rand" | ||
| 6 | "encoding/pem" | ||
| 7 | "fmt" | ||
| 8 | "io" | ||
| 9 | "net" | ||
| 10 | "os" | ||
| 11 | "os/exec" | ||
| 12 | "path/filepath" | ||
| 13 | "strings" | ||
| 14 | "testing" | ||
| 15 | "time" | ||
| 16 | |||
| 17 | "golang.org/x/crypto/ssh" | ||
| 18 | |||
| 19 | "gitbay.org/gitbay/internal/config" | ||
| 20 | "gitbay.org/gitbay/internal/control" | ||
| 21 | "gitbay.org/gitbay/internal/gitutil" | ||
| 22 | "gitbay.org/gitbay/internal/packlimit" | ||
| 23 | "gitbay.org/gitbay/internal/store" | ||
| 24 | ) | ||
| 25 | |||
| 26 | // A client that connects and sends nothing is cut at | ||
| 27 | // push_receive_timeout: the idle rule has not started, since no pack | ||
| 28 | // has begun. Its slot comes back. | ||
| 29 | func TestSilentPushKilledAtReceiveTimeout(t *testing.T) { | ||
| 30 | cfg, st, alice := cloneFixture(t) | ||
| 31 | cfg.Limits.PushIdle = "200ms" | ||
| 32 | cfg.Limits.PushReceiveTimeout = "1s" | ||
| 33 | key := store.SSHKey{ID: 1, Scope: "full", Fingerprint: "SHA256:test"} | ||
| 34 | pushes := packlimit.New(1, 1, 0, time.Second) | ||
| 35 | codec := make(chan int, 1) | ||
| 36 | start := time.Now() | ||
| 37 | go func() { | ||
| 38 | codec <- Exec(cfg, st, nil, pushes, alice, key, control.Term{}, "git-receive-pack alice/app", | ||
| 39 | silentStdin(t), io.Discard, io.Discard, nil, nil, nil) | ||
| 40 | }() | ||
| 41 | slotTaken(t, pushes) | ||
| 42 | pushEnded(t, codec, pushes, 5*time.Second) | ||
| 43 | if d := time.Since(start); d < time.Second { | ||
| 44 | t.Fatalf("killed after %s, before push_receive_timeout", d) | ||
| 45 | } | ||
| 46 | } | ||
| 47 | |||
| 48 | // pktLine frames s as one pkt-line. | ||
| 49 | func pktLine(s string) string { return fmt.Sprintf("%04x%s", len(s)+4, s) } | ||
| 50 | |||
| 51 | const zeroSHA = "0000000000000000000000000000000000000000" | ||
| 52 | |||
| 53 | // Once the pack has begun, a client that goes silent is cut at | ||
| 54 | // push_idle, long before push_receive_timeout. | ||
| 55 | func TestPushKilledAtIdleOncePackStarts(t *testing.T) { | ||
| 56 | cfg, st, alice := cloneFixture(t) | ||
| 57 | cfg.Limits.PushIdle = "300ms" | ||
| 58 | key := store.SSHKey{ID: 1, Scope: "full", Fingerprint: "SHA256:test"} | ||
| 59 | pushes := packlimit.New(1, 1, 0, time.Second) | ||
| 60 | r, w, err := os.Pipe() | ||
| 61 | if err != nil { | ||
| 62 | t.Fatal(err) | ||
| 63 | } | ||
| 64 | t.Cleanup(func() { r.Close(); w.Close() }) | ||
| 65 | io.WriteString(w, pktLine(zeroSHA+" "+strings.Repeat("1", 40)+" refs/heads/main\x00report-status\n")+"0000PA") | ||
| 66 | codec := make(chan int, 1) | ||
| 67 | start := time.Now() | ||
| 68 | go func() { | ||
| 69 | codec <- Exec(cfg, st, nil, pushes, alice, key, control.Term{}, "git-receive-pack alice/app", | ||
| 70 | r, io.Discard, io.Discard, nil, nil, nil) | ||
| 71 | }() | ||
| 72 | slotTaken(t, pushes) | ||
| 73 | // The signature split across two writes still arms the rule. | ||
| 74 | time.Sleep(100 * time.Millisecond) | ||
| 75 | io.WriteString(w, "CK") | ||
| 76 | pushEnded(t, codec, pushes, 5*time.Second) | ||
| 77 | if d := time.Since(start); d > 5*time.Second { | ||
| 78 | t.Fatalf("killed after %s", d) | ||
| 79 | } | ||
| 80 | } | ||
| 81 | |||
| 82 | // A client silent for longer than push_idle between its commands and | ||
| 83 | // its pack, as while pack-objects compresses a large first push, still | ||
| 84 | // completes. | ||
| 85 | func TestPushSilentBeforePackCompletes(t *testing.T) { | ||
| 86 | cfg, st, alice := cloneFixture(t) | ||
| 87 | cfg.Limits.PushIdle = "300ms" | ||
| 88 | key := store.SSHKey{ID: 1, Scope: "full", Fingerprint: "SHA256:test"} | ||
| 89 | pushes := packlimit.New(1, 1, 0, time.Second) | ||
| 90 | |||
| 91 | src := t.TempDir() | ||
| 92 | git := func(stdin string, args ...string) []byte { | ||
| 93 | t.Helper() | ||
| 94 | cmd := exec.Command("git", append([]string{"-C", src, "-c", "user.name=t", "-c", "user.email=t@t"}, args...)...) | ||
| 95 | cmd.Stdin = strings.NewReader(stdin) | ||
| 96 | out, err := cmd.Output() | ||
| 97 | if err != nil { | ||
| 98 | t.Fatalf("git %v: %v", args, err) | ||
| 99 | } | ||
| 100 | return out | ||
| 101 | } | ||
| 102 | git("", "init", "-q", "-b", "main") | ||
| 103 | if err := os.WriteFile(filepath.Join(src, "README"), []byte("x\n"), 0o644); err != nil { | ||
| 104 | t.Fatal(err) | ||
| 105 | } | ||
| 106 | git("", "add", ".") | ||
| 107 | git("", "commit", "-q", "-m", "one") | ||
| 108 | sha := strings.TrimSpace(string(git("", "rev-parse", "HEAD"))) | ||
| 109 | pack := git(sha+"\n", "pack-objects", "--revs", "--stdout", "-q") | ||
| 110 | |||
| 111 | pr, pw := io.Pipe() | ||
| 112 | go func() { | ||
| 113 | io.WriteString(pw, pktLine(zeroSHA+" "+sha+" refs/heads/main\x00report-status\n")+"0000") | ||
| 114 | time.Sleep(time.Second) | ||
| 115 | pw.Write(pack) | ||
| 116 | pw.Close() | ||
| 117 | }() | ||
| 118 | var errOut strings.Builder | ||
| 119 | if code := Exec(cfg, st, nil, pushes, alice, key, control.Term{}, "git-receive-pack alice/app", | ||
| 120 | pr, io.Discard, &errOut, nil, nil, nil); code != 0 { | ||
| 121 | t.Fatalf("exit %d: %s", code, errOut.String()) | ||
| 122 | } | ||
| 123 | out, err := exec.Command("git", "-C", control.RepoDir(cfg.Server.Root, "alice", "app"), "rev-parse", "refs/heads/main").Output() | ||
| 124 | if err != nil || strings.TrimSpace(string(out)) != sha { | ||
| 125 | t.Fatalf("main is %q (%v), want %s", out, err, sha) | ||
| 126 | } | ||
| 127 | } | ||
| 128 | |||
| 129 | // A client that trickles bytes stays clear of push_idle but is cut when | ||
| 130 | // pre-receive has not started push_receive_timeout after the slot. | ||
| 131 | func TestTricklingPushKilledAtReceiveTimeout(t *testing.T) { | ||
| 132 | cfg, st, alice := cloneFixture(t) | ||
| 133 | cfg.Limits.PushIdle = "300ms" | ||
| 134 | cfg.Limits.PushReceiveTimeout = "1s" | ||
| 135 | key := store.SSHKey{ID: 1, Scope: "full", Fingerprint: "SHA256:test"} | ||
| 136 | pushes := packlimit.New(1, 1, 0, time.Second) | ||
| 137 | r, w, err := os.Pipe() | ||
| 138 | if err != nil { | ||
| 139 | t.Fatal(err) | ||
| 140 | } | ||
| 141 | t.Cleanup(func() { r.Close(); w.Close() }) | ||
| 142 | stop := make(chan struct{}) | ||
| 143 | defer close(stop) | ||
| 144 | go func() { | ||
| 145 | // One long pkt-line whose payload arrives a byte at a time. | ||
| 146 | if _, err := io.WriteString(w, "fff0"); err != nil { | ||
| 147 | return | ||
| 148 | } | ||
| 149 | for { | ||
| 150 | select { | ||
| 151 | case <-stop: | ||
| 152 | return | ||
| 153 | case <-time.After(50 * time.Millisecond): | ||
| 154 | if _, err := w.Write([]byte("a")); err != nil { | ||
| 155 | return | ||
| 156 | } | ||
| 157 | } | ||
| 158 | } | ||
| 159 | }() | ||
| 160 | codec := make(chan int, 1) | ||
| 161 | start := time.Now() | ||
| 162 | go func() { | ||
| 163 | codec <- Exec(cfg, st, nil, pushes, alice, key, control.Term{}, "git-receive-pack alice/app", | ||
| 164 | r, io.Discard, io.Discard, nil, nil, nil) | ||
| 165 | }() | ||
| 166 | slotTaken(t, pushes) | ||
| 167 | pushEnded(t, codec, pushes, 5*time.Second) | ||
| 168 | if d := time.Since(start); d < time.Second { | ||
| 169 | t.Fatalf("killed after %s, before push_receive_timeout", d) | ||
| 170 | } | ||
| 171 | } | ||
| 172 | |||
| 173 | // A real push whose post-receive outlasts push_idle completes: git's | ||
| 174 | // side-band keepalives count as bytes to the client. | ||
| 175 | func TestPushSurvivesLongPostReceive(t *testing.T) { | ||
| 176 | if _, err := exec.LookPath("ssh"); err != nil { | ||
| 177 | t.Skip("no ssh client") | ||
| 178 | } | ||
| 179 | root := t.TempDir() | ||
| 180 | st, err := store.Open(filepath.Join(root, "gitbay.db")) | ||
| 181 | if err != nil { | ||
| 182 | t.Fatal(err) | ||
| 183 | } | ||
| 184 | t.Cleanup(func() { st.Close() }) | ||
| 185 | if err := st.MigrateUp(); err != nil { | ||
| 186 | t.Fatal(err) | ||
| 187 | } | ||
| 188 | uid, err := st.CreateUser("alice", false) | ||
| 189 | if err != nil { | ||
| 190 | t.Fatal(err) | ||
| 191 | } | ||
| 192 | if _, err := st.CreateRepo("user", uid, "app", "public"); err != nil { | ||
| 193 | t.Fatal(err) | ||
| 194 | } | ||
| 195 | hooks := t.TempDir() | ||
| 196 | if err := os.WriteFile(filepath.Join(hooks, "post-receive"), []byte("#!/bin/sh\nsleep 4\n"), 0o755); err != nil { | ||
| 197 | t.Fatal(err) | ||
| 198 | } | ||
| 199 | if err := gitutil.InitBare(control.RepoDir(root, "alice", "app"), "main", hooks); err != nil { | ||
| 200 | t.Fatal(err) | ||
| 201 | } | ||
| 202 | |||
| 203 | pub, priv, err := ed25519.GenerateKey(rand.Reader) | ||
| 204 | if err != nil { | ||
| 205 | t.Fatal(err) | ||
| 206 | } | ||
| 207 | sshPub, err := ssh.NewPublicKey(pub) | ||
| 208 | if err != nil { | ||
| 209 | t.Fatal(err) | ||
| 210 | } | ||
| 211 | if err := st.AddSSHKey(uid, ssh.FingerprintSHA256(sshPub), sshPub.Type(), sshPub.Marshal(), "full", "test"); err != nil { | ||
| 212 | t.Fatal(err) | ||
| 213 | } | ||
| 214 | block, err := ssh.MarshalPrivateKey(priv, "") | ||
| 215 | if err != nil { | ||
| 216 | t.Fatal(err) | ||
| 217 | } | ||
| 218 | keyFile := filepath.Join(t.TempDir(), "id_ed25519") | ||
| 219 | if err := os.WriteFile(keyFile, pem.EncodeToMemory(block), 0o600); err != nil { | ||
| 220 | t.Fatal(err) | ||
| 221 | } | ||
| 222 | |||
| 223 | cfg := config.Default() | ||
| 224 | cfg.Server.Root = root | ||
| 225 | cfg.Limits.PushIdle = "2s" | ||
| 226 | pushes := packlimit.New(1, 1, 0, time.Second) | ||
| 227 | srv, err := New(cfg, st, nil, pushes) | ||
| 228 | if err != nil { | ||
| 229 | t.Fatal(err) | ||
| 230 | } | ||
| 231 | ln, err := net.Listen("tcp", "127.0.0.1:0") | ||
| 232 | if err != nil { | ||
| 233 | t.Fatal(err) | ||
| 234 | } | ||
| 235 | go srv.Serve(ln) | ||
| 236 | t.Cleanup(func() { ln.Close() }) | ||
| 237 | port := ln.Addr().(*net.TCPAddr).Port | ||
| 238 | |||
| 239 | src := t.TempDir() | ||
| 240 | git := func(args ...string) string { | ||
| 241 | t.Helper() | ||
| 242 | cmd := exec.Command("git", args...) | ||
| 243 | cmd.Dir = src | ||
| 244 | cmd.Env = append(os.Environ(), | ||
| 245 | "GIT_CONFIG_GLOBAL=/dev/null", "GIT_CONFIG_NOSYSTEM=1", | ||
| 246 | "GIT_AUTHOR_NAME=a", "GIT_AUTHOR_EMAIL=a@example.test", | ||
| 247 | "GIT_COMMITTER_NAME=a", "GIT_COMMITTER_EMAIL=a@example.test", | ||
| 248 | fmt.Sprintf("GIT_SSH_COMMAND=ssh -F /dev/null -i %s -o IdentitiesOnly=yes -o StrictHostKeyChecking=no -o UserKnownHostsFile=/dev/null -o LogLevel=ERROR -p %d", keyFile, port)) | ||
| 249 | out, err := cmd.CombinedOutput() | ||
| 250 | if err != nil { | ||
| 251 | t.Fatalf("git %v: %v\n%s", args, err, out) | ||
| 252 | } | ||
| 253 | return string(out) | ||
| 254 | } | ||
| 255 | git("init", "-q", "-b", "main") | ||
| 256 | if err := os.WriteFile(filepath.Join(src, "README"), []byte("x\n"), 0o644); err != nil { | ||
| 257 | t.Fatal(err) | ||
| 258 | } | ||
| 259 | git("add", ".") | ||
| 260 | git("commit", "-q", "-m", "one") | ||
| 261 | start := time.Now() | ||
| 262 | git("push", "git@127.0.0.1:alice/app", "main") | ||
| 263 | if d := time.Since(start); d < 4*time.Second { | ||
| 264 | t.Fatalf("push took %s; the post-receive did not run", d) | ||
| 265 | } | ||
| 266 | out, err := exec.Command("git", "-C", control.RepoDir(root, "alice", "app"), "rev-parse", "refs/heads/main").CombinedOutput() | ||
| 267 | if err != nil || strings.TrimSpace(string(out)) == "" { | ||
| 268 | t.Fatalf("main not pushed: %v %s", err, out) | ||
| 269 | } | ||
| 270 | r, err := pushes.Acquire(nil, "elsewhere") | ||
| 271 | if err != nil { | ||
| 272 | t.Fatalf("slot not released: %v", err) | ||
| 273 | } | ||
| 274 | r() | ||
| 275 | } | ||
internal/sshd/refusal_test.go +30 −5
| @@ -345,26 +345,51 @@ func TestPushPerPrincipalCap(t *testing.T) { | |||
| 345 | } | 345 | } |
| 346 | } | 346 | } |
| 347 | 347 | ||
| 348 | // A revoked key kills a push waiting on its client, and the slot comes | 348 | // A revoked key kills a push waiting on its client, and the slot it |
| 349 | // back. | 349 | // held comes back. |
| 350 | func TestPushKilledWhenKeyRevokedReleasesSlot(t *testing.T) { | 350 | func TestPushKilledWhenKeyRevokedReleasesSlot(t *testing.T) { |
| 351 | cfg, st, alice := cloneFixture(t) | 351 | cfg, st, alice := cloneFixture(t) |
| 352 | key := store.SSHKey{ID: 1, Scope: "full", Fingerprint: "SHA256:test"} | 352 | key := store.SSHKey{ID: 1, Scope: "full", Fingerprint: "SHA256:test"} |
| 353 | pushes := packlimit.New(1, 1, 0, time.Second) | 353 | pushes := packlimit.New(1, 1, 0, time.Second) |
| 354 | revoked := make(chan struct{}) | ||
| 354 | codec := make(chan int, 1) | 355 | codec := make(chan int, 1) |
| 355 | go func() { | 356 | go func() { |
| 356 | codec <- Exec(cfg, st, nil, pushes, alice, key, control.Term{}, "git-receive-pack alice/app", | 357 | codec <- Exec(cfg, st, nil, pushes, alice, key, control.Term{}, "git-receive-pack alice/app", |
| 357 | silentStdin(t), io.Discard, io.Discard, nil, nil, closed()) | 358 | silentStdin(t), io.Discard, io.Discard, nil, nil, revoked) |
| 358 | }() | 359 | }() |
| 360 | slotTaken(t, pushes) | ||
| 361 | close(revoked) | ||
| 362 | pushEnded(t, codec, pushes, 5*time.Second) | ||
| 363 | } | ||
| 364 | |||
| 365 | // slotTaken waits until l's one slot is held. | ||
| 366 | func slotTaken(t *testing.T, l *packlimit.Limiter) { | ||
| 367 | t.Helper() | ||
| 368 | deadline := time.Now().Add(5 * time.Second) | ||
| 369 | for time.Now().Before(deadline) { | ||
| 370 | r, err := l.Acquire(closed(), "probe") | ||
| 371 | if err != nil { | ||
| 372 | return | ||
| 373 | } | ||
| 374 | r() | ||
| 375 | time.Sleep(10 * time.Millisecond) | ||
| 376 | } | ||
| 377 | t.Fatal("the push never took the slot") | ||
| 378 | } | ||
| 379 | |||
| 380 | // pushEnded requires a push to end, killed, within within, with l's | ||
| 381 | // slot free again. | ||
| 382 | func pushEnded(t *testing.T, codec <-chan int, l *packlimit.Limiter, within time.Duration) { | ||
| 383 | t.Helper() | ||
| 359 | select { | 384 | select { |
| 360 | case code := <-codec: | 385 | case code := <-codec: |
| 361 | if code != protocol.ExitFailure { | 386 | if code != protocol.ExitFailure { |
| 362 | t.Fatalf("exit %d, want the push killed", code) | 387 | t.Fatalf("exit %d, want the push killed", code) |
| 363 | } | 388 | } |
| 364 | case <-time.After(5 * time.Second): | 389 | case <-time.After(within): |
| 365 | t.Fatal("push still running") | 390 | t.Fatal("push still running") |
| 366 | } | 391 | } |
| 367 | r, err := pushes.Acquire(nil, "elsewhere") | 392 | r, err := l.Acquire(nil, "elsewhere") |
| 368 | if err != nil { | 393 | if err != nil { |
| 369 | t.Fatalf("slot not released after the kill: %v", err) | 394 | t.Fatalf("slot not released after the kill: %v", err) |
| 370 | } | 395 | } |
internal/sshd/sshd.go +91 −2
| @@ -3,6 +3,7 @@ | |||
| 3 | package sshd | 3 | package sshd |
| 4 | 4 | ||
| 5 | import ( | 5 | import ( |
| 6 | "bytes" | ||
| 6 | "context" | 7 | "context" |
| 7 | "crypto/ed25519" | 8 | "crypto/ed25519" |
| 8 | "crypto/rand" | 9 | "crypto/rand" |
| @@ -529,6 +530,7 @@ func runGit(cfg config.Config, st *store.Store, packs, pushes *packlimit.Limiter | |||
| 529 | return protocol.ExitUsage | 530 | return protocol.ExitUsage |
| 530 | } | 531 | } |
| 531 | write := service == "git-receive-pack" | 532 | write := service == "git-receive-pack" |
| 533 | cancel, keepAlive := revoked, 0 | ||
| 532 | 534 | ||
| 533 | repo, err := st.RepoByPath(argv[1]) | 535 | repo, err := st.RepoByPath(argv[1]) |
| 534 | if err != nil { | 536 | if err != nil { |
| @@ -612,6 +614,7 @@ func runGit(cfg config.Config, st *store.Store, packs, pushes *packlimit.Limiter | |||
| 612 | if code != protocol.ExitOK { | 614 | if code != protocol.ExitOK { |
| 613 | return code | 615 | return code |
| 614 | } | 616 | } |
| 617 | slotAt := time.Now() | ||
| 615 | // Deferred before Transport runs, so it fires after receive-pack | 618 | // Deferred before Transport runs, so it fires after receive-pack |
| 616 | // has exited, on every path: the client hanging up ends its | 619 | // has exited, on every path: the client hanging up ends its |
| 617 | // stdin and receive-pack with it, and a revoked key kills it. | 620 | // stdin and receive-pack with it, and a revoked key kills it. |
| @@ -624,8 +627,62 @@ func runGit(cfg config.Config, st *store.Store, packs, pushes *packlimit.Limiter | |||
| 624 | } | 627 | } |
| 625 | defer st.DeletePushToken(token) | 628 | defer st.DeletePushToken(token) |
| 626 | env = append(env, hookd.EnvToken+"="+token) | 629 | env = append(env, hookd.EnvToken+"="+token) |
| 630 | idleFor, receiveFor := cfg.Limits.PushTimeouts() | ||
| 631 | keepAlive = pushKeepAlive(idleFor) | ||
| 632 | if pushes != nil { | ||
| 633 | // A client that holds its slot while sending nothing, or | ||
| 634 | // trickles its pack, is cut: when pre-receive has not | ||
| 635 | // started receiveFor after the slot was taken, or after | ||
| 636 | // idleFor with no byte either way once the pack has begun | ||
| 637 | // or pre-receive has started, whichever is first. | ||
| 638 | // receive-pack's keepalives count, so indexing and hooks do | ||
| 639 | // not end it. The idle rule waits for the pack because the | ||
| 640 | // client sends nothing while pack-objects counts and | ||
| 641 | // compresses, and receive-pack sends no keepalive then. | ||
| 642 | // Once pre-receive starts only the idle rule applies, so | ||
| 643 | // post-receive is never cut short by the clock. | ||
| 644 | client, clientOut := stdin, stdout | ||
| 645 | var idle <-chan struct{} | ||
| 646 | var arm, unwatch func() | ||
| 647 | stdin, stdout, idle, arm, unwatch = packlimit.Idle(stdin, stdout, idleFor) | ||
| 648 | stdin = &packStart{r: stdin, seen: arm} | ||
| 649 | defer unwatch() | ||
| 650 | started, forget := hookd.AwaitPreReceive(token) | ||
| 651 | defer forget() | ||
| 652 | deadline := time.NewTimer(time.Until(slotAt.Add(receiveFor))) | ||
| 653 | defer deadline.Stop() | ||
| 654 | kill := make(chan struct{}) | ||
| 655 | finished := make(chan struct{}) | ||
| 656 | defer close(finished) | ||
| 657 | go func() { | ||
| 658 | receiving := deadline.C | ||
| 659 | for { | ||
| 660 | select { | ||
| 661 | case <-finished: | ||
| 662 | return | ||
| 663 | case <-started: | ||
| 664 | arm() | ||
| 665 | started, receiving = nil, nil | ||
| 666 | continue | ||
| 667 | case <-revoked: | ||
| 668 | case <-idle: | ||
| 669 | case <-receiving: | ||
| 670 | } | ||
| 671 | close(kill) | ||
| 672 | // A read blocked on a silent client outlives git; | ||
| 673 | // closing the channel ends it and the stdin copy, so | ||
| 674 | // Transport's Wait returns. | ||
| 675 | for _, c := range []any{client, clientOut} { | ||
| 676 | if c, ok := c.(io.Closer); ok { | ||
| 677 | c.Close() | ||
| 678 | } | ||
| 679 | } | ||
| 680 | return | ||
| 681 | } | ||
| 682 | }() | ||
| 683 | cancel = kill | ||
| 684 | } | ||
| 627 | } | 685 | } |
| 628 | cancel := revoked | ||
| 629 | if !write { | 686 | if !write { |
| 630 | // Pack generation shares one budget with smart HTTP and git://. | 687 | // Pack generation shares one budget with smart HTTP and git://. |
| 631 | release, code := takeSlot(packs, "user:"+strconv.FormatInt(user.ID, 10), done, stderr, | 688 | release, code := takeSlot(packs, "user:"+strconv.FormatInt(user.ID, 10), done, stderr, |
| @@ -679,7 +736,7 @@ func runGit(cfg config.Config, st *store.Store, packs, pushes *packlimit.Limiter | |||
| 679 | }() | 736 | }() |
| 680 | cancel = kill | 737 | cancel = kill |
| 681 | } | 738 | } |
| 682 | if err := gitutil.Transport(service, dir, stdin, stdout, stderr, env, maxPack, cancel); err != nil { | 739 | if err := gitutil.Transport(service, dir, stdin, stdout, stderr, env, maxPack, keepAlive, cancel); err != nil { |
| 683 | return protocol.ExitFailure | 740 | return protocol.ExitFailure |
| 684 | } | 741 | } |
| 685 | return protocol.ExitOK | 742 | return protocol.ExitOK |
| @@ -702,3 +759,35 @@ func takeSlot(l *packlimit.Limiter, principal string, done <-chan struct{}, stde | |||
| 702 | } | 759 | } |
| 703 | return nil, protocol.ExitFailure | 760 | return nil, protocol.ExitFailure |
| 704 | } | 761 | } |
| 762 | |||
| 763 | // pushKeepAlive is receive.keepAlive, in seconds, for a push idle limit | ||
| 764 | // of idle: several keepalives fit in one idle period, and never more | ||
| 765 | // than five seconds apart. | ||
| 766 | func pushKeepAlive(idle time.Duration) int { | ||
| 767 | return max(1, min(5, int(idle/(4*time.Second)))) | ||
| 768 | } | ||
| 769 | |||
| 770 | // packStart calls seen once the pack signature "PACK" has passed | ||
| 771 | // through r, which may split it across reads. It can fire early on a | ||
| 772 | // command or push option that holds the word; that only starts the idle | ||
| 773 | // rule sooner. | ||
| 774 | type packStart struct { | ||
| 775 | r io.Reader | ||
| 776 | seen func() | ||
| 777 | tail []byte // up to three bytes carried from the previous read | ||
| 778 | done bool | ||
| 779 | } | ||
| 780 | |||
| 781 | func (p *packStart) Read(b []byte) (int, error) { | ||
| 782 | n, err := p.r.Read(b) | ||
| 783 | if n > 0 && !p.done { | ||
| 784 | buf := append(p.tail, b[:n]...) | ||
| 785 | if bytes.Contains(buf, []byte("PACK")) { | ||
| 786 | p.done, p.tail = true, nil | ||
| 787 | p.seen() | ||
| 788 | } else { | ||
| 789 | p.tail = append([]byte(nil), buf[max(0, len(buf)-3):]...) | ||
| 790 | } | ||
| 791 | } | ||
| 792 | return n, err | ||
| 793 | } | ||