limits: bound concurrent pushes; repo download under the pack limit !542

merged merged by cmc on 2026-09-29 17:11 UTC · krz/gitbay:push-limit-308 into main

29 files changed, +1235 −69

Layout: unified · split

.gitbay/wiki/Admin.org +40 −3
@@ -357,8 +357,8 @@ push=.
357 repositories may take; a push may be no larger than what is left. 357 repositories may take; a push may be no larger than what is left.
358- =pack_concurrency= (3), =pack_per_principal= (2), =pack_queue= (32), 358- =pack_concurrency= (3), =pack_per_principal= (2), =pack_queue= (32),
359 =pack_queue_wait= (="60s"=) — git pack generation (clones, fetches, 359 =pack_queue_wait= (="60s"=) — git pack generation (clones, fetches,
360 =git archive --remote=, web archive downloads) over SSH, smart HTTP 360 =git archive --remote=, web archive downloads, =repo download= over
361 and git:// shares one 361 SSH and the API) over SSH, smart HTTP and git:// shares one
362 budget: this many at once, this many per account (per client 362 budget: this many at once, this many per account (per client
363 address when anonymous: an IPv4 address, or an IPv6 /64), and this 363 address when anonymous: an IPv4 address, or an IPv6 /64), and this
364 many waiting for at most the wait. Anonymous clients together hold 364 many waiting for at most the wait. Anonymous clients together hold
@@ -375,11 +375,48 @@ push=.
375 for two minutes: a client reading below about 550 B/s, or an HTTP 375 for two minutes: a client reading below about 550 B/s, or an HTTP
376 request body that takes over two minutes with nothing written back, 376 request body that takes over two minutes with nothing written back,
377 is cut. Ref listings (info/refs, protocol v2 377 is cut. Ref listings (info/refs, protocol v2
378 =ls-refs=), pushes and =repo download= are outside the budget. For the 378 =ls-refs=) and pushes are outside the budget. A busy =repo download=
379 exits 1 with the same message; over the API it is a 503 with
380 =Retry-After: 60=. For the
379 three counts 0 means the default and a negative value turns that 381 three counts 0 means the default and a negative value turns that
380 bound off. The defaults suit a four-core host; see [[Performance]]. 382 bound off. The defaults suit a four-core host; see [[Performance]].
381 With =ssh.mode = "system"= each SSH session is its own process and 383 With =ssh.mode = "system"= each SSH session is its own process and
382 SSH clones are not counted. 384 SSH clones are not counted.
385- =push_concurrency= (2), =push_per_principal= (1), =push_queue= (16),
386 =push_queue_wait= (="60s"=) — =git receive-pack= over SSH, on a
387 budget of its own so clones cannot starve pushes or the reverse:
388 this many at once, this many per account or deploy key, and this
389 many waiting for at most the wait. One principal may have up to half
390 of =push_queue= (at least one) waiting, so parallel pushes from one
391 account queue rather than being refused, but cannot fill the queue.
392 A deploy key is its own principal, apart from the account that
393 registered it, so an account with deploy keys can hold
394 =push_per_principal= slots per key;
395 the global cap still holds. A bot account such as a runner's counts
396 like any other account. The slot is taken after the access checks and held
397 until receive-pack exits, which is after =post-receive= (merge
398 detection, CI queueing, webhooks) has run. Past that the client gets
399 "the server is busy: it is at its limit of concurrent pushes…" and
400 exit 1, and the daemon logs a =push limit= warning at most once a
401 minute. A client that disconnects while queued leaves the queue; one
402 that disconnects mid-push ends receive-pack and frees its slot, and
403 a revoked key kills it. 0 and negative values read as for =pack_*=.
404 HTTP and git:// carry no pushes. With =ssh.mode = "system"= pushes
405 are not counted and the two timeouts below do not apply.
406- =push_idle= (="60s"=), =push_receive_timeout= (="15m"=) — a push
407 holding a slot is killed, its process group with it, when its
408 pre-receive has not started =push_receive_timeout= after it took the
409 slot, or when no byte has been read from the client and none written
410 to it for =push_idle=. The idle rule starts when the pack signature
411 arrives from the client or pre-receive starts, whichever is first
412 (a delete-only push sends no pack): before that the client may be
413 silent while its =pack-objects= counts and compresses, and only
414 =push_receive_timeout= applies. receive-pack runs with
415 =receive.keepAlive= set to a quarter of =push_idle= (1 to 5 seconds),
416 so the keepalives it sends to side-band clients while it indexes and
417 runs hooks keep a working push alive. Once pre-receive has started
418 only =push_idle= applies: post-receive is never cut short by the
419 clock. A push killed before pre-receive answers updates no refs.
383- =max_pack_bytes= (2 GiB) — the largest pack one push may send, 420- =max_pack_bytes= (2 GiB) — the largest pack one push may send,
384 enforced as =receive.maxInputSize= and lowered to what an owner's 421 enforced as =receive.maxInputSize= and lowered to what an owner's
385 storage quota has left. 422 storage quota has left.
.gitbay/wiki/Architecture/09-Controls.org +2 −1
@@ -97,7 +97,8 @@ chapter names of OWASP ASVS 4.0 where one fits.
97| Control | Status | Evidence | 97| Control | Status | Evidence |
98|---------------------------------------------+----------+------------------------------------------------------------------| 98|---------------------------------------------+----------+------------------------------------------------------------------|
99| Rate limits on API and writes | in place | [[file:05-Identity-and-Access.org][5. Rate limits]] | 99| Rate limits on API and writes | in place | [[file:05-Identity-and-Access.org][5. Rate limits]] |
100| Concurrency limit on git pack generation | in place | global, per-principal, bounded queue across SSH, HTTP and git:// (=internal/packlimit=); not in system SSH mode | 100| Concurrency limit on git pack generation | in place | global, per-principal, bounded queue across SSH, HTTP and git://, =repo download= included (=internal/packlimit=); not in system SSH mode |
101| Concurrency limit on pushes | in place | receive-pack on its own global, per-principal, bounded-queue budget; a deploy key is its own principal; killed at =push_receive_timeout= before pre-receive, or after =push_idle= with nothing moving once the pack has begun (=internal/sshd/sshd.go=, =internal/packlimit=, =internal/hookd=); not in system SSH mode |
101| Service hardening | in place | systemd sandboxing ([[file:03-Deployment.org][3]]) | 102| Service hardening | in place | systemd sandboxing ([[file:03-Deployment.org][3]]) |
102| Backups offsite and append-only | in place | restic with append-only credentials (documented) | 103| Backups offsite and append-only | in place | restic with append-only credentials (documented) |
103| Restore tested | partial | drill 2026-09-29 from the offsite copy (Admin wiki "Restore drill"); secrets not checked in that drill; the off-host =secret.key= exists (#305) and is checked in the next | 104| Restore tested | partial | drill 2026-09-29 from the offsite copy (Admin wiki "Restore drill"); secrets not checked in that drill; the off-host =secret.key= exists (#305) and is checked in the next |
.gitbay/wiki/Architecture/10-Known-Gaps.org +2 −4
@@ -16,10 +16,8 @@ what the 2026-09-27 review found; remove a row when its issue closes.
16| Area | Gap | Severity | 16| Area | Gap | Severity |
17|-------+-------------------------------------------------------------------------------------------------------------+----------| 17|-------+-------------------------------------------------------------------------------------------------------------+----------|
18| Audit | The hash chain is unkeyed, so whoever can write the database can edit a row and recompute every later hash; removing the newest audit rows, or writing new rows under their freed ids, needs no recomputing at all. Neither is detectable from the database; only comparing =gitbayd admin audit verify='s last id and hash with the daemon's journal shows it. Rows written by =gitbayd shell= (=ssh.mode = "system"=) and host admin commands have no journal copy, and the refusal caps are per process, so under that mode each connection counts separately | low | 18| Audit | The hash chain is unkeyed, so whoever can write the database can edit a row and recompute every later hash; removing the newest audit rows, or writing new rows under their freed ids, needs no recomputing at all. Neither is detectable from the database; only comparing =gitbayd admin audit verify='s last id and hash with the daemon's journal shows it. Rows written by =gitbayd shell= (=ssh.mode = "system"=) and host admin commands have no journal copy, and the refusal caps are per process, so under that mode each connection counts separately | low |
19| Access | Grants and parked profile about texts (=profile_about_backfill=) of deleted accounts and organizations, and deploy keys of deleted repositories, left by deletes before #306, are removed on upgrade and the "schema migrated" log line gives the counts. One whose id a later row had already taken is no longer an orphan and stays with that row; on gitbay.org the orphans found (two grants, one deploy key) name ids no later row had taken | low | 19| Availability | Under =ssh.mode = "system"= each SSH session is a separate =gitbayd shell= process, so the pack and push limits (=internal/packlimit=, #262, #308) cannot count SSH clones or pushes across sessions; only HTTP and git:// share a budget there | low |
20| Availability | Under =ssh.mode = "system"= each SSH session is a separate =gitbayd shell= process, so the pack-generation limit (=internal/packlimit=, #262) cannot count SSH clones across sessions; only HTTP and git:// share a budget there | low | 20| Availability | A push holds its slot for at most =push_receive_timeout= (15m) before its pre-receive starts, and, once its pack has begun, at most =push_idle= (60s) with nothing moving either way. A client that trickles its pack and reconnects when cut can hold a slot per principal for up to the receive deadline each time; deploy keys multiply the per-principal share, bounded by =push_concurrency= | low |
21| Availability | Pushes have no concurrency limit; =max_pack_bytes= bounds each one, not how many run at once | medium |
22| Availability | =repo download= (SSH, API) runs =git archive= outside the pack limit; only its two-minute deadline and 512 MiB cap bound it | low |
23| Availability | An HTTP or git:// client that disconnects while queued for a pack slot keeps its place until =pack_queue_wait= runs out; only SSH notices the disconnect | low | 21| Availability | An HTTP or git:// client that disconnects while queued for a pack slot keeps its place until =pack_queue_wait= runs out; only SSH notices the disconnect | low |
24 22
25* Questions an auditor will ask that have no answer yet 23* Questions an auditor will ask that have no answer yet
.gitbay/wiki/Performance.org +5 −1
@@ -47,7 +47,11 @@ The practical ceiling on this hardware is concurrent pack generation:
47full clones of large repositories are CPU-bound in git itself (the 17s 47full clones of large repositories are CPU-bound in git itself (the 17s
48clone ran git at ~156% CPU). =limits.pack_concurrency= bounds how many 48clone ran git at ~156% CPU). =limits.pack_concurrency= bounds how many
49run at once across SSH, HTTP and git://, with a queue behind it (see 49run at once across SSH, HTTP and git://, with a queue behind it (see
50[[Admin]], =[limits]=). 50[[Admin]], =[limits]=). Pushes are the other git-bound load: receive-pack
51indexes the pack it is sent, about a core for a large one, and then runs
52its hooks. =limits.push_concurrency= (2) bounds those on a separate
53budget, one per account or deploy key by default, so a clone storm and
54a push storm each leave the other its slots.
51 55
52* Concurrent clones 56* Concurrent clones
53 57
.gitbay/wiki/Threat-Model.org +4 −3
@@ -367,9 +367,10 @@ is recorded here rather than in a closed issue:
367- *Timing and traffic analysis.* Token comparison is a hash index lookup 367- *Timing and traffic analysis.* Token comparison is a hash index lookup
368 by design, but nothing has been measured. 368 by design, but nothing has been measured.
369- *Denial of service by resource exhaustion* beyond rate. Concurrent 369- *Denial of service by resource exhaustion* beyond rate. Concurrent
370 clones, fetches and web archives are bounded by the pack limit 370 clones, fetches, web archives and =repo download= are bounded by the
371 (#262); pushes are not, beyond =max_pack_bytes= on each one. Nor are 371 pack limit (#262, #308), and pushes by a separate push limit (#308)
372 pathological diffs, deep histories, or zip bombs in LFS. 372 on top of =max_pack_bytes= on each one. Pathological diffs, deep
373 histories and zip bombs in LFS are not.
373 374
374A sweep is a point in time. This section says what a reader should not 375A sweep is a point in time. This section says what a reader should not
375assume has been checked. 376assume has been checked.
CHANGELOG.org +22
@@ -4,6 +4,28 @@ Versioning follows semver from v0.1.0. Database migrations run
4automatically on daemon start; upgrade notes appear per release when 4automatically on daemon start; upgrade notes appear per release when
5anything beyond "replace the binary and restart" is needed. 5anything beyond "replace the binary and restart" is needed.
6 6
7* Unreleased
8
9- Pushes have a concurrency limit of their own, separate from the pack
10 budget: =[limits] push_concurrency= (2), =push_per_principal= (1),
11 =push_queue= (16) and =push_queue_wait= ("60s"), read like the
12 =pack_*= settings. One principal may have up to half of =push_queue=
13 waiting, so parallel pushes from one account queue. A deploy key
14 counts apart from the account that registered it. The slot is held until receive-pack and its
15 post-receive hook have finished; a busy push exits 1 with "the server
16 is busy: it is at its limit of concurrent pushes…", and the daemon
17 logs a =push limit= warning at most once a minute. (#308)
18- A push holding a slot is killed when its pre-receive has not started
19 =[limits] push_receive_timeout= ("15m") after it took the slot, or,
20 once its pack has begun or pre-receive has started, when nothing has
21 moved either way for =push_idle= ("60s").
22 receive-pack runs with =receive.keepAlive= set from =push_idle=, so
23 long hooks do not trip it. (#308)
24- =repo download= over SSH and the API takes a slot from the pack
25 limit, released once =git archive= exits; busy is exit 1 with the
26 clone message, and 503 with =Retry-After: 60= on both API
27 endpoints. (#308)
28
7* v1.40.0 — 2026-09-29 29* v1.40.0 — 2026-09-29
8 30
9DKIM verification for reply mail, and ids that are not reused after a 31DKIM verification for reply mail, and ids that are not reused after a
cmd/gitbayd/main.go +15 −1
@@ -262,12 +262,20 @@ func serveCmd() *cobra.Command {
262 if packMax > 1 { 262 if packMax > 1 {
263 packs.CapClass("ip:", packMax-1) 263 packs.CapClass("ip:", packMax-1)
264 } 264 }
265 // receive-pack has a budget of its own. Parallel pushes
266 // from one account (scripts, bots, several terminals)
267 // queue, up to half the queue, rather than being refused
268 // once one is waiting.
269 pushMax, pushPer, pushQueue, pushWait := cfg.Limits.PushLimits()
270 pushes := packlimit.New(pushMax, pushPer, pushQueue, pushWait)
271 pushes.Name("push")
272 pushes.CapQueue(pushPerQueue(pushQueue))
265 273
266 errCh := make(chan error, 3) 274 errCh := make(chan error, 3)
267 var sshSrv *sshd.Server 275 var sshSrv *sshd.Server
268 var sshLn, gitLn net.Listener 276 var sshLn, gitLn net.Listener
269 if cfg.SSH.Mode == "embedded" { 277 if cfg.SSH.Mode == "embedded" {
270 srv, err := sshd.New(cfg, st, packs) 278 srv, err := sshd.New(cfg, st, packs, pushes)
271 if err != nil { 279 if err != nil {
272 return err 280 return err
273 } 281 }
@@ -645,3 +653,9 @@ func reapPending(ctx context.Context, st *store.Store, maxAge time.Duration) {
645 } 653 }
646 } 654 }
647} 655}
656
657// pushPerQueue is how many pushes one principal may have waiting: half
658// the queue, at least one.
659func pushPerQueue(queue int) int {
660 return max(1, queue/2)
661}
cmd/gitbayd/system.go +2 −2
@@ -100,8 +100,8 @@ func shellCmd() *cobra.Command {
100 os.Exit(protocol.ExitUsage) 100 os.Exit(protocol.ExitUsage)
101 } 101 }
102 // Each forced command is its own process, so there is no 102 // Each forced command is its own process, so there is no
103 // shared pack budget in system mode. 103 // shared pack or push budget in system mode.
104 code := sshd.Exec(cfg, st, nil, user, key, control.ParseTerm(os.Getenv("GITBAY_TERM")), cmdline, os.Stdin, os.Stdout, os.Stderr, nil, nil, nil) 104 code := sshd.Exec(cfg, st, nil, nil, user, key, control.ParseTerm(os.Getenv("GITBAY_TERM")), cmdline, os.Stdin, os.Stdout, os.Stderr, nil, nil, nil)
105 st.Close() 105 st.Close()
106 os.Exit(code) 106 os.Exit(code)
107 return nil 107 return nil
internal/config/config.go +72 −9
@@ -35,6 +35,21 @@ const (
35 DefaultPackQueueWait = time.Minute 35 DefaultPackQueueWait = time.Minute
36) 36)
37 37
38// Push defaults: receive-pack indexes what it is sent, a core each for a
39// large push, and has its own budget so clones and pushes cannot starve
40// each other.
41const (
42 DefaultPushConcurrency = 2
43 DefaultPushPerPrincipal = 1
44 DefaultPushQueue = 16
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
51)
52
38type Config struct { 53type Config struct {
39 Server Server `toml:"server"` 54 Server Server `toml:"server"`
40 SSH SSH `toml:"ssh"` 55 SSH SSH `toml:"ssh"`
@@ -235,12 +250,53 @@ type Limits struct {
235 PackPerPrincipal int `toml:"pack_per_principal"` 250 PackPerPrincipal int `toml:"pack_per_principal"`
236 PackQueue int `toml:"pack_queue"` 251 PackQueue int `toml:"pack_queue"`
237 PackQueueWait string `toml:"pack_queue_wait"` 252 PackQueueWait string `toml:"pack_queue_wait"`
253 // PushConcurrency caps receive-pack running at once over SSH, with
254 // its hooks; PushPerPrincipal caps it per account or deploy key.
255 // PushQueue and PushQueueWait bound the wait for a slot. The counts
256 // read like the pack_* ones: 0 takes the default, negative is off.
257 PushConcurrency int `toml:"push_concurrency"`
258 PushPerPrincipal int `toml:"push_per_principal"`
259 PushQueue int `toml:"push_queue"`
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"`
238} 270}
239 271
240// PackLimits resolves the pack_* settings for packlimit.New. A zero 272// PackLimits resolves the pack_* settings for packlimit.New. A zero
241// max or per is no bound; an unbounded queue is math.MaxInt, since 273// max or per is no bound; an unbounded queue is math.MaxInt, since
242// packlimit reads a zero queue as no queue at all. 274// packlimit reads a zero queue as no queue at all.
243func (l Limits) PackLimits() (max, per, queue int, wait time.Duration) { 275func (l Limits) PackLimits() (max, per, queue int, wait time.Duration) {
276 return resolveLimits(l.PackConcurrency, l.PackPerPrincipal, l.PackQueue, l.PackQueueWait,
277 DefaultPackConcurrency, DefaultPackPerPrincipal, DefaultPackQueue, DefaultPackQueueWait)
278}
279
280// PushLimits resolves the push_* settings the way PackLimits does the
281// pack_* ones.
282func (l Limits) PushLimits() (max, per, queue int, wait time.Duration) {
283 return resolveLimits(l.PushConcurrency, l.PushPerPrincipal, l.PushQueue, l.PushQueueWait,
284 DefaultPushConcurrency, DefaultPushPerPrincipal, DefaultPushQueue, DefaultPushQueueWait)
285}
286
287// PushTimeouts resolves push_idle and push_receive_timeout.
288func (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
299func resolveLimits(maxV, perV, queueV int, waitV string, maxDef, perDef, queueDef int, waitDef time.Duration) (max, per, queue int, wait time.Duration) {
244 pick := func(v, def int) int { 300 pick := func(v, def int) int {
245 switch { 301 switch {
246 case v == 0: 302 case v == 0:
@@ -250,16 +306,15 @@ func (l Limits) PackLimits() (max, per, queue int, wait time.Duration) {
250 } 306 }
251 return v 307 return v
252 } 308 }
253 queue = pick(l.PackQueue, DefaultPackQueue) 309 queue = pick(queueV, queueDef)
254 if l.PackQueue < 0 { 310 if queueV < 0 {
255 queue = math.MaxInt 311 queue = math.MaxInt
256 } 312 }
257 wait = DefaultPackQueueWait 313 wait = waitDef
258 if d, err := time.ParseDuration(l.PackQueueWait); err == nil && d > 0 { 314 if d, err := time.ParseDuration(waitV); err == nil && d > 0 {
259 wait = d 315 wait = d
260 } 316 }
261 return pick(l.PackConcurrency, DefaultPackConcurrency), 317 return pick(maxV, maxDef), pick(perV, perDef), queue, wait
262 pick(l.PackPerPrincipal, DefaultPackPerPrincipal), queue, wait
263} 318}
264 319
265type Mail struct { 320type Mail struct {
@@ -619,9 +674,17 @@ func (c Config) Validate() error {
619 if c.Limits.MaxReposPerUser < 0 || c.Limits.MaxBytesPerUser < 0 || c.Limits.MaxSnippetsPerUser < 0 { 674 if c.Limits.MaxReposPerUser < 0 || c.Limits.MaxBytesPerUser < 0 || c.Limits.MaxSnippetsPerUser < 0 {
620 errs = append(errs, errors.New("limits.max_repos_per_user, max_bytes_per_user and max_snippets_per_user must not be negative")) 675 errs = append(errs, errors.New("limits.max_repos_per_user, max_bytes_per_user and max_snippets_per_user must not be negative"))
621 } 676 }
622 if w := c.Limits.PackQueueWait; w != "" { 677 for _, w := range []struct{ name, val string }{
623 if d, err := time.ParseDuration(w); err != nil || d <= 0 { 678 {"pack_queue_wait", c.Limits.PackQueueWait},
624 errs = append(errs, fmt.Errorf("limits.pack_queue_wait %q must be a positive duration such as 60s", w)) 679 {"push_queue_wait", c.Limits.PushQueueWait},
680 {"push_idle", c.Limits.PushIdle},
681 {"push_receive_timeout", c.Limits.PushReceiveTimeout},
682 } {
683 if w.val == "" {
684 continue
685 }
686 if d, err := time.ParseDuration(w.val); err != nil || d <= 0 {
687 errs = append(errs, fmt.Errorf("limits.%s %q must be a positive duration such as 60s", w.name, w.val))
625 } 688 }
626 } 689 }
627 if c.Push.Enabled { 690 if c.Push.Enabled {
internal/config/config_test.go +48
@@ -63,6 +63,39 @@ func TestPackLimits(t *testing.T) {
63 } 63 }
64} 64}
65 65
66// push_* resolve like pack_*, from their own defaults.
67func TestPushLimits(t *testing.T) {
68 max, per, queue, wait := Limits{}.PushLimits()
69 if max != DefaultPushConcurrency || per != DefaultPushPerPrincipal || queue != DefaultPushQueue || wait != DefaultPushQueueWait {
70 t.Fatalf("defaults: %d %d %d %s", max, per, queue, wait)
71 }
72 if max != 2 || per != 1 || queue != 16 || wait != time.Minute {
73 t.Fatalf("defaults moved: %d %d %d %s", max, per, queue, wait)
74 }
75 max, per, queue, wait = Limits{PushConcurrency: -1, PushPerPrincipal: -1, PushQueue: -1, PushQueueWait: "5s"}.PushLimits()
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)
78 }
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 {
81 t.Fatal(err)
82 }
83 max, per, queue, wait = cfg.Limits.PushLimits()
84 if max != 4 || per != 2 || queue != 8 || wait != 30*time.Second {
85 t.Fatalf("loaded: %d %d %d %s", max, per, queue, wait)
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 }
93 // The pack budget is not read from the push settings.
94 if pm, _, _, _ := cfg.Limits.PackLimits(); pm != DefaultPackConcurrency {
95 t.Fatalf("pack_concurrency %d, want the default", pm)
96 }
97}
98
66func TestContradictions(t *testing.T) { 99func TestContradictions(t *testing.T) {
67 cases := []struct { 100 cases := []struct {
68 name string 101 name string
@@ -74,6 +107,21 @@ func TestContradictions(t *testing.T) {
74 minimal + "\n[limits]\npack_queue_wait = \"soon\"\n", 107 minimal + "\n[limits]\npack_queue_wait = \"soon\"\n",
75 "limits.pack_queue_wait", 108 "limits.pack_queue_wait",
76 }, 109 },
110 {
111 "bad push_queue_wait",
112 minimal + "\n[limits]\npush_queue_wait = \"-5s\"\n",
113 "limits.push_queue_wait",
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 },
77 { 125 {
78 "registration open without smtp", 126 "registration open without smtp",
79 minimal + "\n[registration]\nmode = \"open\"\n", 127 minimal + "\n[registration]\nmode = \"open\"\n",
internal/control/control.go +7
@@ -14,6 +14,7 @@ import (
14 "time" 14 "time"
15 15
16 "gitbay.org/gitbay/internal/config" 16 "gitbay.org/gitbay/internal/config"
17 "gitbay.org/gitbay/internal/packlimit"
17 "gitbay.org/gitbay/internal/protocol" 18 "gitbay.org/gitbay/internal/protocol"
18 "gitbay.org/gitbay/internal/store" 19 "gitbay.org/gitbay/internal/store"
19) 20)
@@ -67,6 +68,12 @@ type Ctx struct {
67 // restarting. It closes Done too; a command that ends on Done checks 68 // restarting. It closes Done too; a command that ends on Done checks
68 // it to say why. 69 // it to say why.
69 Stopping <-chan struct{} 70 Stopping <-chan struct{}
71 // Packs is the pack-generation limiter a command that runs git to
72 // produce an archive takes a slot from; nil is no limit.
73 Packs *packlimit.Limiter
74 // Busy is set when a limiter turned the command away, so the API
75 // can answer 503 with Retry-After rather than a failure.
76 Busy bool
70} 77}
71 78
72// SourceWeb is Ctx.Source for a request from a browser session. Its 79// SourceWeb is Ctx.Source for a request from a browser session. Its
internal/control/download_test.go added +71
@@ -0,0 +1,71 @@
1package control
2
3import (
4 "bytes"
5 "strings"
6 "testing"
7 "time"
8
9 "gitbay.org/gitbay/internal/packlimit"
10 "gitbay.org/gitbay/internal/protocol"
11 "gitbay.org/gitbay/internal/store"
12)
13
14// repo download takes a pack slot: busy while the limiter is full, and
15// the slot is free again once git archive has exited.
16func TestRepoDownloadTakesAPackSlot(t *testing.T) {
17 st, repo, root, _ := prunedRepo(t)
18 alice, err := st.UserByUsername("alice")
19 if err != nil {
20 t.Fatal(err)
21 }
22 packs := packlimit.New(1, 0, 0, time.Second)
23 hold, err := packs.Acquire(nil, "ip:elsewhere")
24 if err != nil {
25 t.Fatal(err)
26 }
27 c, errOut := pruneCtx(st, root, alice)
28 c.Packs = packs
29 if code := Dispatch(c, []string{"repo", "download", repo.Path()}); code != protocol.ExitFailure ||
30 !strings.Contains(errOut.String(), "limit of concurrent clones") {
31 t.Fatalf("busy: exit %d: %q", code, errOut.String())
32 }
33 hold()
34
35 c, errOut = pruneCtx(st, root, alice)
36 c.Packs = packs
37 if code := Dispatch(c, []string{"repo", "download", repo.Path()}); code != protocol.ExitOK {
38 t.Fatalf("free: exit %d: %s", code, errOut.String())
39 }
40 if c.Stdout.(*bytes.Buffer).Len() == 0 {
41 t.Fatal("no archive written")
42 }
43 hold, err = packs.Acquire(nil, "ip:elsewhere")
44 if err != nil {
45 t.Fatalf("slot not released after the download: %v", err)
46 }
47 hold()
48}
49
50// A download the caller may not make is refused before the limiter.
51func TestRepoDownloadRefusalStaysOffLimiter(t *testing.T) {
52 st, repo, root, _ := prunedRepo(t)
53 if err := st.SetRepoVisibility(repo.ID, "private"); err != nil {
54 t.Fatal(err)
55 }
56 bobID, err := st.CreateUser("bob", false)
57 if err != nil {
58 t.Fatal(err)
59 }
60 packs := packlimit.New(1, 0, 0, time.Second)
61 hold, err := packs.Acquire(nil, "ip:elsewhere")
62 if err != nil {
63 t.Fatal(err)
64 }
65 defer hold()
66 c, errOut := pruneCtx(st, root, store.User{ID: bobID, Username: "bob"})
67 c.Packs = packs
68 if code := Dispatch(c, []string{"repo", "download", repo.Path()}); code != protocol.ExitNotFound {
69 t.Fatalf("exit %d: %q", code, errOut.String())
70 }
71}
internal/control/explore.go +20
@@ -1,10 +1,13 @@
1package control 1package control
2 2
3import ( 3import (
4 "errors"
4 "io" 5 "io"
6 "strconv"
5 "strings" 7 "strings"
6 8
7 "gitbay.org/gitbay/internal/gitutil" 9 "gitbay.org/gitbay/internal/gitutil"
10 "gitbay.org/gitbay/internal/packlimit"
8 "gitbay.org/gitbay/internal/policy" 11 "gitbay.org/gitbay/internal/policy"
9 "gitbay.org/gitbay/internal/protocol" 12 "gitbay.org/gitbay/internal/protocol"
10) 13)
@@ -117,6 +120,23 @@ func runRepoDownload(c *Ctx, args []string) int {
117 // The prefix git puts on every path inside the archive, so unpacking 120 // The prefix git puts on every path inside the archive, so unpacking
118 // lands in a named directory rather than the current one. 121 // lands in a named directory rather than the current one.
119 prefix := repo.Name + "-" + ref 122 prefix := repo.Name + "-" + ref
123 // git archive draws on the same budget as clones and web archives,
124 // counted against the account the way upload-pack is.
125 principal := "user:" + strconv.FormatInt(c.User.ID, 10)
126 release, err := c.Packs.Acquire(c.Done, principal)
127 if err != nil {
128 transport := "ssh"
129 if c.ViaAPI {
130 transport = "api"
131 }
132 c.Packs.Refused(transport, principal, err)
133 c.Busy = true
134 if errors.Is(err, packlimit.ErrBusy) {
135 return c.fail(protocol.ExitFailure, "the server is busy: it is at its limit of concurrent clones and fetches; try again in a minute")
136 }
137 return c.fail(protocol.ExitFailure, "the server is restarting; try again in a minute")
138 }
139 defer release()
120 if err := gitutil.Archive(dir, ref, prefix, c.Stdout); err != nil { 140 if err := gitutil.Archive(dir, ref, prefix, c.Stdout); err != nil {
121 return c.fail(protocol.ExitFailure, "%v", err) 141 return c.fail(protocol.ExitFailure, "%v", err)
122 } 142 }
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.
41func Transport(service, repoPath string, stdin io.Reader, stdout, errW io.Writer, extraEnv []string, maxPack int64, cancel <-chan struct{}) error { 42func 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.
333var (
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.
342func 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
354func 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.
205func 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/httpd/api.go +9
@@ -77,6 +77,7 @@ func (s *Server) apiCmd(w http.ResponseWriter, r *http.Request) {
77 Expires: tok.ExpiresAt, 77 Expires: tok.ExpiresAt,
78 Done: s.until(r), 78 Done: s.until(r),
79 Stopping: s.stopping, 79 Stopping: s.stopping,
80 Packs: s.packs,
80 } 81 }
81 code := control.Dispatch(ctx, req.Argv) 82 code := control.Dispatch(ctx, req.Argv)
82 83
@@ -96,12 +97,20 @@ func (s *Server) apiCmd(w http.ResponseWriter, r *http.Request) {
96 body["stderr"] = msg 97 body["stderr"] = msg
97 } 98 }
98 w.Header().Set("Content-Type", "application/json") 99 w.Header().Set("Content-Type", "application/json")
100 if ctx.Busy {
101 status = http.StatusServiceUnavailable
102 w.Header().Set("Retry-After", busyRetryAfter)
103 }
99 w.WriteHeader(status) 104 w.WriteHeader(status)
100 json.NewEncoder(w).Encode(body) 105 json.NewEncoder(w).Encode(body)
101} 106}
102 107
103// statusForExit maps a command's exit code onto an HTTP status, shared by 108// statusForExit maps a command's exit code onto an HTTP status, shared by
104// both API surfaces so they cannot answer the same failure differently. 109// both API surfaces so they cannot answer the same failure differently.
110// busyRetryAfter is the Retry-After on a 503 for a command a limiter
111// turned away.
112const busyRetryAfter = "60"
113
105func statusForExit(code int) int { 114func statusForExit(code int) int {
106 switch code { 115 switch code {
107 case protocol.ExitOK: 116 case protocol.ExitOK:
internal/httpd/apibusy_test.go added +52
@@ -0,0 +1,52 @@
1package httpd
2
3import (
4 "net/http"
5 "net/http/httptest"
6 "strings"
7 "testing"
8
9 "gitbay.org/gitbay/internal/control"
10 "gitbay.org/gitbay/internal/gitutil"
11 "gitbay.org/gitbay/internal/store"
12)
13
14// repo download turned away by a full pack limit is 503 with
15// Retry-After on both API endpoints, not a 500.
16func TestAPIBusyDownloadIs503(t *testing.T) {
17 s := busyServer(t)
18 s.apiLimit = newAPILimiter(0)
19 if err := gitutil.InitBare(control.RepoDir(s.cfg.Server.Root, "alice", "app"), "main", t.TempDir()); err != nil {
20 t.Fatal(err)
21 }
22 seed(t, s, 10)
23 alice, err := s.st.UserByUsername("alice")
24 if err != nil {
25 t.Fatal(err)
26 }
27 if err := s.st.CreateAPIToken(alice.ID, "t", store.HashToken("secret"), "full", nil, 0); err != nil {
28 t.Fatal(err)
29 }
30 check := func(what string, w *httptest.ResponseRecorder) {
31 t.Helper()
32 if w.Code != http.StatusServiceUnavailable || w.Header().Get("Retry-After") != "60" ||
33 !strings.Contains(w.Body.String(), "busy") {
34 t.Fatalf("%s: status %d, Retry-After %q, body %s", what, w.Code, w.Header().Get("Retry-After"), w.Body.String())
35 }
36 if w.Header().Get("ETag") != "" {
37 t.Fatalf("%s: a busy answer carries an ETag", what)
38 }
39 }
40
41 r := httptest.NewRequest("POST", "/api/v1/cmd", strings.NewReader(`{"argv":["repo","download","alice/app"]}`))
42 r.Header.Set("Authorization", "Bearer secret")
43 w := httptest.NewRecorder()
44 s.apiCmd(w, r)
45 check("POST", w)
46
47 r = httptest.NewRequest("GET", "/api/v1/read?argv=repo&argv=download&argv=alice/app", nil)
48 r.Header.Set("Authorization", "Bearer secret")
49 w = httptest.NewRecorder()
50 s.apiRead(w, r)
51 check("GET", w)
52}
internal/httpd/apiread.go +10
@@ -64,6 +64,7 @@ func (s *Server) apiRead(w http.ResponseWriter, r *http.Request) {
64 ReadOnly: true, 64 ReadOnly: true,
65 Done: s.until(r), 65 Done: s.until(r),
66 Stopping: s.stopping, 66 Stopping: s.stopping,
67 Packs: s.packs,
67 } 68 }
68 code := control.Dispatch(ctx, argv) 69 code := control.Dispatch(ctx, argv)
69 70
@@ -81,6 +82,15 @@ func (s *Server) apiRead(w http.ResponseWriter, r *http.Request) {
81 return 82 return
82 } 83 }
83 84
85 if ctx.Busy {
86 // Not cacheable: the next try may succeed.
87 w.Header().Set("Content-Type", "application/json")
88 w.Header().Set("Retry-After", busyRetryAfter)
89 w.WriteHeader(http.StatusServiceUnavailable)
90 w.Write(payload)
91 return
92 }
93
84 // Responses are authorized per account, so the ETag is salted with the 94 // Responses are authorized per account, so the ETag is salted with the
85 // caller: two users asking the same question may get different answers, 95 // caller: two users asking the same question may get different answers,
86 // and neither should ever be served the other's. 96 // and neither should ever be served the other's.
internal/packlimit/packlimit.go +30 −8
@@ -1,8 +1,9 @@
1// Package packlimit bounds concurrent git pack generation. upload-pack 1// Package packlimit bounds concurrent git processes. upload-pack and
2// and upload-archive over SSH, smart HTTP and git:// draw on one 2// upload-archive over SSH, smart HTTP and git:// draw on one budget,
3// budget: a global cap, a cap per principal (an account, or a client 3// receive-pack on a second: each a global cap, a cap per principal (an
4// address on the anonymous transports), and a bounded queue whose 4// account, a deploy key, or a client address on the anonymous
5// waiters give up after a fixed wait or when the client goes away. 5// transports), and a bounded queue whose waiters give up after a fixed
6// wait or when the client goes away.
6// Waiters are not served in order; a new arrival can take a freed slot 7// Waiters are not served in order; a new arrival can take a freed slot
7// ahead of them, and the wait bounds how long any one of them waits. 8// ahead of them, and the wait bounds how long any one of them waits.
8package packlimit 9package packlimit
@@ -23,7 +24,9 @@ var (
23 24
24type Limiter struct { 25type Limiter struct {
25 max, per, queue int 26 max, per, queue int
27 perQueue int // waiting, per principal; per unless set
26 wait time.Duration 28 wait time.Duration
29 name string // what is limited, for the refusal log
27 30
28 // Principals starting with class may hold at most classCap slots 31 // Principals starting with class may hold at most classCap slots
29 // between them; classCap 0 is no class cap. 32 // between them; classCap 0 is no class cap.
@@ -45,7 +48,7 @@ func New(max, per, queue int, wait time.Duration) *Limiter {
45 if max <= 0 { 48 if max <= 0 {
46 return nil 49 return nil
47 } 50 }
48 return &Limiter{max: max, per: per, queue: queue, wait: wait, 51 return &Limiter{max: max, per: per, perQueue: per, queue: queue, wait: wait, name: "pack",
49 held: map[string]int{}, waiting: map[string]int{}, changed: make(chan struct{}), 52 held: map[string]int{}, waiting: map[string]int{}, changed: make(chan struct{}),
50 warned: map[string]time.Time{}} 53 warned: map[string]time.Time{}}
51} 54}
@@ -71,10 +74,29 @@ func (l *Limiter) Refused(transport, principal string, err error) {
71 if errors.Is(err, ErrGone) { 74 if errors.Is(err, ErrGone) {
72 reason = "gone" 75 reason = "gone"
73 } 76 }
74 slog.Warn("pack limit: request turned away (logged at most once a minute per transport)", 77 slog.Warn(l.name+" limit: request turned away (logged at most once a minute per transport)",
75 "transport", transport, "class", class, "reason", reason) 78 "transport", transport, "class", class, "reason", reason)
76} 79}
77 80
81// Name sets what the refusal log calls this limit ("pack" unless set).
82// Call it before the limiter is in use.
83func (l *Limiter) Name(name string) {
84 if l == nil {
85 return
86 }
87 l.name = name
88}
89
90// CapQueue lets one principal have up to n requests waiting, where by
91// default it may have as many as it may run. It applies only while a
92// per-principal cap is set. Call it before the limiter is in use.
93func (l *Limiter) CapQueue(n int) {
94 if l == nil {
95 return
96 }
97 l.perQueue = n
98}
99
78// CapClass caps the slots that principals starting with prefix may hold 100// CapClass caps the slots that principals starting with prefix may hold
79// between them. Call it before the limiter is in use. 101// between them. Call it before the limiter is in use.
80func (l *Limiter) CapClass(prefix string, n int) { 102func (l *Limiter) CapClass(prefix string, n int) {
@@ -116,7 +138,7 @@ func (l *Limiter) Acquire(done <-chan struct{}, principal string) (release func(
116 l.mu.Unlock() 138 l.mu.Unlock()
117 return l.releaser(principal), nil 139 return l.releaser(principal), nil
118 } 140 }
119 if l.queued >= l.queue || (l.per > 0 && l.waiting[principal] >= l.per) { 141 if l.queued >= l.queue || (l.per > 0 && l.waiting[principal] >= l.perQueue) {
120 l.mu.Unlock() 142 l.mu.Unlock()
121 return nil, ErrBusy 143 return nil, ErrBusy
122 } 144 }
internal/packlimit/packlimit_test.go +52
@@ -313,3 +313,55 @@ func TestRefusedLogsOncePerTransport(t *testing.T) {
313 var none *Limiter 313 var none *Limiter
314 none.Refused("git", "ip:x", ErrBusy) 314 none.Refused("git", "ip:x", ErrBusy)
315} 315}
316
317// Two limiters log apart: a named one says what it limits, and its
318// once-a-minute window does not silence the other's.
319func TestRefusedNamesTheLimit(t *testing.T) {
320 var buf bytes.Buffer
321 old := slog.Default()
322 slog.SetDefault(slog.New(slog.NewTextHandler(&buf, nil)))
323 t.Cleanup(func() { slog.SetDefault(old) })
324
325 packs := New(1, 0, 0, time.Second)
326 pushes := New(1, 0, 0, time.Second)
327 pushes.Name("push")
328 packs.Refused("ssh", "user:4", ErrBusy)
329 pushes.Refused("ssh", "key:9", ErrBusy)
330 out := buf.String()
331 if !strings.Contains(out, `"pack limit: request turned away`) || !strings.Contains(out, `"push limit: request turned away`) {
332 t.Fatalf("want one line per limit:\n%s", out)
333 }
334 if !strings.Contains(out, "class=key") || strings.Contains(out, "key:9") {
335 t.Fatalf("deploy key principal logged or class missing:\n%s", out)
336 }
337 var none *Limiter
338 none.Name("push")
339}
340
341// With a waiting cap above per, one principal runs per and queues up to
342// perQueue, and is refused past that.
343func TestPerPrincipalQueueCap(t *testing.T) {
344 l := New(4, 1, 16, 5*time.Second)
345 l.CapQueue(4)
346 r, err := l.Acquire(nil, "a")
347 if err != nil {
348 t.Fatal(err)
349 }
350 for i := 1; i <= 4; i++ {
351 go func() {
352 if r, err := l.Acquire(nil, "a"); err == nil {
353 r()
354 }
355 }()
356 waitQueued(t, l, i)
357 }
358 if _, err := l.Acquire(nil, "a"); !errors.Is(err, ErrBusy) {
359 t.Fatalf("sixth for a: %v", err)
360 }
361 rb, err := l.Acquire(nil, "b")
362 if err != nil {
363 t.Fatalf("b blocked by a's queue: %v", err)
364 }
365 rb()
366 r()
367}
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.
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}
internal/packlimit/watch_test.go +35
@@ -2,6 +2,7 @@ package packlimit
2 2
3import ( 3import (
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.
45func 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 @@
1package sshd
2
3import (
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.
29func 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.
49func pktLine(s string) string { return fmt.Sprintf("%04x%s", len(s)+4, s) }
50
51const 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.
55func 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.
85func 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.
131func 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.
175func 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 +181 −7
@@ -2,9 +2,11 @@ package sshd
2 2
3import ( 3import (
4 "bytes" 4 "bytes"
5 "errors"
5 "io" 6 "io"
6 "os" 7 "os"
7 "path/filepath" 8 "path/filepath"
9 "strconv"
8 "strings" 10 "strings"
9 "testing" 11 "testing"
10 "time" 12 "time"
@@ -56,7 +58,7 @@ func TestRefusedPushIsAudited(t *testing.T) {
56 bob.Pending = pending 58 bob.Pending = pending
57 key := store.SSHKey{Scope: "full", Fingerprint: "SHA256:test"} 59 key := store.SSHKey{Scope: "full", Fingerprint: "SHA256:test"}
58 var out, errOut bytes.Buffer 60 var out, errOut bytes.Buffer
59 code := Exec(cfg, st, nil, bob, key, control.Term{}, "git-receive-pack alice/app", 61 code := Exec(cfg, st, nil, nil, bob, key, control.Term{}, "git-receive-pack alice/app",
60 strings.NewReader(""), &out, &errOut, nil, nil, nil) 62 strings.NewReader(""), &out, &errOut, nil, nil, nil)
61 if code != protocol.ExitDenied { 63 if code != protocol.ExitDenied {
62 t.Fatalf("pending %v: exit %d: %s", pending, code, errOut.String()) 64 t.Fatalf("pending %v: exit %d: %s", pending, code, errOut.String())
@@ -83,7 +85,7 @@ func TestCloneRefusedWhenPackSlotsAreFull(t *testing.T) {
83 key := store.SSHKey{Scope: "full", Fingerprint: "SHA256:test"} 85 key := store.SSHKey{Scope: "full", Fingerprint: "SHA256:test"}
84 for _, service := range []string{"git-upload-pack", "git-upload-archive"} { 86 for _, service := range []string{"git-upload-pack", "git-upload-archive"} {
85 var out, errOut bytes.Buffer 87 var out, errOut bytes.Buffer
86 code := Exec(cfg, st, packs, bob, key, control.Term{}, service+" alice/app", 88 code := Exec(cfg, st, packs, nil, bob, key, control.Term{}, service+" alice/app",
87 strings.NewReader(""), &out, &errOut, nil, nil, nil) 89 strings.NewReader(""), &out, &errOut, nil, nil, nil)
88 if code != protocol.ExitFailure || !strings.Contains(errOut.String(), "busy") { 90 if code != protocol.ExitFailure || !strings.Contains(errOut.String(), "busy") {
89 t.Fatalf("%s: exit %d: %q", service, code, errOut.String()) 91 t.Fatalf("%s: exit %d: %q", service, code, errOut.String())
@@ -106,7 +108,7 @@ func TestPushBypassesPackLimitAndCloneReleasesSlot(t *testing.T) {
106 packs := packlimit.New(1, 0, 0, time.Second) 108 packs := packlimit.New(1, 0, 0, time.Second)
107 109
108 var out, errOut bytes.Buffer 110 var out, errOut bytes.Buffer
109 if code := Exec(cfg, st, packs, alice, key, control.Term{}, "git-upload-pack alice/app", 111 if code := Exec(cfg, st, packs, nil, alice, key, control.Term{}, "git-upload-pack alice/app",
110 strings.NewReader("0000"), &out, &errOut, nil, nil, nil); code != protocol.ExitOK { 112 strings.NewReader("0000"), &out, &errOut, nil, nil, nil); code != protocol.ExitOK {
111 t.Fatalf("clone: exit %d: %s", code, errOut.String()) 113 t.Fatalf("clone: exit %d: %s", code, errOut.String())
112 } 114 }
@@ -118,7 +120,7 @@ func TestPushBypassesPackLimitAndCloneReleasesSlot(t *testing.T) {
118 120
119 out.Reset() 121 out.Reset()
120 errOut.Reset() 122 errOut.Reset()
121 if code := Exec(cfg, st, packs, alice, key, control.Term{}, "git-receive-pack alice/app", 123 if code := Exec(cfg, st, packs, nil, alice, key, control.Term{}, "git-receive-pack alice/app",
122 strings.NewReader("0000"), &out, &errOut, nil, nil, nil); code != protocol.ExitOK { 124 strings.NewReader("0000"), &out, &errOut, nil, nil, nil); code != protocol.ExitOK {
123 t.Fatalf("push with slots full: exit %d: %s", code, errOut.String()) 125 t.Fatalf("push with slots full: exit %d: %s", code, errOut.String())
124 } 126 }
@@ -159,7 +161,7 @@ func killedClone(t *testing.T, stdout io.Writer, done, stopping, revoked <-chan
159 packs := packlimit.New(1, 0, 0, time.Second) 161 packs := packlimit.New(1, 0, 0, time.Second)
160 codec := make(chan int, 1) 162 codec := make(chan int, 1)
161 go func() { 163 go func() {
162 codec <- Exec(cfg, st, packs, alice, key, control.Term{}, "git-upload-pack alice/app", 164 codec <- Exec(cfg, st, packs, nil, alice, key, control.Term{}, "git-upload-pack alice/app",
163 silentStdin(t), stdout, io.Discard, done, stopping, revoked) 165 silentStdin(t), stdout, io.Discard, done, stopping, revoked)
164 }() 166 }()
165 select { 167 select {
@@ -208,7 +210,7 @@ func TestCloneRunsOnDuringRestart(t *testing.T) {
208 key := store.SSHKey{Scope: "full", Fingerprint: "SHA256:test"} 210 key := store.SSHKey{Scope: "full", Fingerprint: "SHA256:test"}
209 packs := packlimit.New(1, 0, 0, time.Second) 211 packs := packlimit.New(1, 0, 0, time.Second)
210 var errOut bytes.Buffer 212 var errOut bytes.Buffer
211 if code := Exec(cfg, st, packs, alice, key, control.Term{}, "git-upload-pack alice/app", 213 if code := Exec(cfg, st, packs, nil, alice, key, control.Term{}, "git-upload-pack alice/app",
212 strings.NewReader("0000"), io.Discard, &errOut, closed(), closed(), nil); code != protocol.ExitOK { 214 strings.NewReader("0000"), io.Discard, &errOut, closed(), closed(), nil); code != protocol.ExitOK {
213 t.Fatalf("exit %d: %s", code, errOut.String()) 215 t.Fatalf("exit %d: %s", code, errOut.String())
214 } 216 }
@@ -233,9 +235,181 @@ func TestRefusedCloneStaysOffLimiter(t *testing.T) {
233 defer hold() 235 defer hold()
234 key := store.SSHKey{Scope: "full", Fingerprint: "SHA256:test"} 236 key := store.SSHKey{Scope: "full", Fingerprint: "SHA256:test"}
235 var out, errOut bytes.Buffer 237 var out, errOut bytes.Buffer
236 code := Exec(cfg, st, packs, bob, key, control.Term{}, "git-upload-pack alice/secret", 238 code := Exec(cfg, st, packs, nil, bob, key, control.Term{}, "git-upload-pack alice/secret",
237 strings.NewReader(""), &out, &errOut, nil, nil, nil) 239 strings.NewReader(""), &out, &errOut, nil, nil, nil)
238 if code != protocol.ExitNotFound || strings.Contains(errOut.String(), "busy") { 240 if code != protocol.ExitNotFound || strings.Contains(errOut.String(), "busy") {
239 t.Fatalf("exit %d: %q", code, errOut.String()) 241 t.Fatalf("exit %d: %q", code, errOut.String())
240 } 242 }
241} 243}
244
245// hungUpStdin is a client that sends nothing until hangUp, which ends
246// its stdin the way a closed channel does.
247func hungUpStdin(t *testing.T) (r *os.File, hangUp func()) {
248 t.Helper()
249 r, w, err := os.Pipe()
250 if err != nil {
251 t.Fatal(err)
252 }
253 t.Cleanup(func() { r.Close(); w.Close() })
254 return r, func() { w.Close() }
255}
256
257// queuedFor waits until principal has a push waiting on l: a probe that
258// gives up at once is then refused busy rather than queued.
259func queuedFor(t *testing.T, l *packlimit.Limiter, principal string) {
260 t.Helper()
261 deadline := time.Now().Add(5 * time.Second)
262 for time.Now().Before(deadline) {
263 if _, err := l.Acquire(closed(), principal); errors.Is(err, packlimit.ErrBusy) {
264 return
265 }
266 time.Sleep(10 * time.Millisecond)
267 }
268 t.Fatalf("no push queued for %s", principal)
269}
270
271// With push_per_principal 1 one account runs one push, queues a second
272// and is refused a third, while another principal still gets in. A
273// client hanging up frees its slot for the one queued behind it.
274func TestPushPerPrincipalCap(t *testing.T) {
275 cfg, st, alice := cloneFixture(t)
276 repo, err := st.RepoByPath("alice/app")
277 if err != nil {
278 t.Fatal(err)
279 }
280 key := store.SSHKey{ID: 1, Scope: "full", Fingerprint: "SHA256:test"}
281 pushes := packlimit.New(2, 1, 16, 10*time.Second)
282 push := func(k store.SSHKey, stdin io.Reader, errOut io.Writer) <-chan int {
283 codec := make(chan int, 1)
284 go func() {
285 codec <- Exec(cfg, st, nil, pushes, alice, k, control.Term{}, "git-receive-pack alice/app",
286 stdin, io.Discard, errOut, nil, nil, nil)
287 }()
288 return codec
289 }
290 exited := func(codec <-chan int, want int, what string) {
291 t.Helper()
292 select {
293 case code := <-codec:
294 if code != want {
295 t.Fatalf("%s: exit %d", what, code)
296 }
297 case <-time.After(5 * time.Second):
298 t.Fatalf("%s still running", what)
299 }
300 }
301 principal := "user:" + strconv.FormatInt(alice.ID, 10)
302
303 in1, hangUp1 := hungUpStdin(t)
304 first := push(key, in1, io.Discard)
305 // The first holds alice's one slot once a probe cannot take it.
306 deadline := time.Now().Add(5 * time.Second)
307 for {
308 r, err := pushes.Acquire(closed(), principal)
309 if err != nil {
310 break
311 }
312 r()
313 if time.Now().After(deadline) {
314 t.Fatal("first push never took a slot")
315 }
316 time.Sleep(10 * time.Millisecond)
317 }
318 in2, hangUp2 := hungUpStdin(t)
319 second := push(key, in2, io.Discard)
320 queuedFor(t, pushes, principal)
321
322 var errOut bytes.Buffer
323 if code := Exec(cfg, st, nil, pushes, alice, key, control.Term{}, "git-receive-pack alice/app",
324 strings.NewReader(""), io.Discard, &errOut, nil, nil, nil); code != protocol.ExitFailure ||
325 !strings.Contains(errOut.String(), "limit of concurrent pushes") {
326 t.Fatalf("third push: exit %d: %q", code, errOut.String())
327 }
328
329 // A deploy key on the same account is its own principal: it takes
330 // the second global slot while alice's push waits.
331 deploy := store.SSHKey{ID: 2, Scope: "deploy:" + strconv.FormatInt(repo.ID, 10) + ":rw", Fingerprint: "SHA256:deploy"}
332 exited(push(deploy, strings.NewReader("0000"), io.Discard), protocol.ExitOK, "deploy key push")
333
334 // receive-pack fails a client that hangs up before sending anything.
335 hangUp1()
336 exited(first, protocol.ExitFailure, "first push")
337 hangUp2()
338 exited(second, protocol.ExitFailure, "second push")
339 for _, p := range []string{"a", "b"} {
340 r, err := pushes.Acquire(nil, p)
341 if err != nil {
342 t.Fatalf("slot not released: %v", err)
343 }
344 defer r()
345 }
346}
347
348// A revoked key kills a push waiting on its client, and the slot it
349// held comes back.
350func TestPushKilledWhenKeyRevokedReleasesSlot(t *testing.T) {
351 cfg, st, alice := cloneFixture(t)
352 key := store.SSHKey{ID: 1, Scope: "full", Fingerprint: "SHA256:test"}
353 pushes := packlimit.New(1, 1, 0, time.Second)
354 revoked := make(chan struct{})
355 codec := make(chan int, 1)
356 go func() {
357 codec <- Exec(cfg, st, nil, pushes, alice, key, control.Term{}, "git-receive-pack alice/app",
358 silentStdin(t), io.Discard, io.Discard, nil, nil, revoked)
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.
366func 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.
382func pushEnded(t *testing.T, codec <-chan int, l *packlimit.Limiter, within time.Duration) {
383 t.Helper()
384 select {
385 case code := <-codec:
386 if code != protocol.ExitFailure {
387 t.Fatalf("exit %d, want the push killed", code)
388 }
389 case <-time.After(within):
390 t.Fatal("push still running")
391 }
392 r, err := l.Acquire(nil, "elsewhere")
393 if err != nil {
394 t.Fatalf("slot not released after the kill: %v", err)
395 }
396 r()
397}
398
399// A push the key may not make is refused before it reaches the limiter.
400func TestRefusedPushStaysOffLimiter(t *testing.T) {
401 cfg, st, bob := execFixture(t)
402 pushes := packlimit.New(1, 0, 0, time.Second)
403 hold, err := pushes.Acquire(nil, "elsewhere")
404 if err != nil {
405 t.Fatal(err)
406 }
407 defer hold()
408 key := store.SSHKey{Scope: "full", Fingerprint: "SHA256:test"}
409 var errOut bytes.Buffer
410 code := Exec(cfg, st, nil, pushes, bob, key, control.Term{}, "git-receive-pack alice/app",
411 strings.NewReader(""), io.Discard, &errOut, nil, nil, nil)
412 if code != protocol.ExitDenied || strings.Contains(errOut.String(), "busy") {
413 t.Fatalf("exit %d: %q", code, errOut.String())
414 }
415}
internal/sshd/sshd.go +142 −24
@@ -3,6 +3,7 @@
3package sshd 3package sshd
4 4
5import ( 5import (
6 "bytes"
6 "context" 7 "context"
7 "crypto/ed25519" 8 "crypto/ed25519"
8 "crypto/rand" 9 "crypto/rand"
@@ -39,6 +40,7 @@ type Server struct {
39 cfg config.Config 40 cfg config.Config
40 st *store.Store 41 st *store.Store
41 packs *packlimit.Limiter 42 packs *packlimit.Limiter
43 pushes *packlimit.Limiter
42 sshCfg *ssh.ServerConfig 44 sshCfg *ssh.ServerConfig
43 authLimiter *rateLimiter 45 authLimiter *rateLimiter
44 sessions sync.WaitGroup // accepted connections still being served 46 sessions sync.WaitGroup // accepted connections still being served
@@ -70,8 +72,8 @@ func (c *conn) cut() {
70 c.net.Close() 72 c.net.Close()
71} 73}
72 74
73func New(cfg config.Config, st *store.Store, packs *packlimit.Limiter) (*Server, error) { 75func New(cfg config.Config, st *store.Store, packs, pushes *packlimit.Limiter) (*Server, error) {
74 s := &Server{cfg: cfg, st: st, packs: packs, authLimiter: newRateLimiter(cfg.Limits.SSHAuthRate, time.Minute), conns: map[*conn]struct{}{}, stopping: make(chan struct{})} 76 s := &Server{cfg: cfg, st: st, packs: packs, pushes: pushes, authLimiter: newRateLimiter(cfg.Limits.SSHAuthRate, time.Minute), conns: map[*conn]struct{}{}, stopping: make(chan struct{})}
75 77
76 sc := &ssh.ServerConfig{ 78 sc := &ssh.ServerConfig{
77 PublicKeyCallback: s.authenticate, 79 PublicKeyCallback: s.authenticate,
@@ -425,7 +427,7 @@ func (s *Server) runExec(c *conn, sconn *ssh.ServerConn, ch ssh.Channel, term co
425 return protocol.ExitDenied 427 return protocol.ExitDenied
426 } 428 }
427 _ = s.st.TouchSSHKey(keyID) 429 _ = s.st.TouchSSHKey(keyID)
428 return Exec(s.cfg, s.st, s.packs, user, key, term, cmdline, ch, ch, ch.Stderr(), done, s.stopping, c.revoked) 430 return Exec(s.cfg, s.st, s.packs, s.pushes, user, key, term, cmdline, ch, ch, ch.Stderr(), done, s.stopping, c.revoked)
429} 431}
430 432
431// runAnonymous handles a session from an unregistered key: the register 433// runAnonymous handles a session from an unregistered key: the register
@@ -461,7 +463,9 @@ func (s *Server) runAnonymous(ch ssh.Channel, keyB64, cmdline string) int {
461// Exec runs one SSH exec command line for an authenticated key. It is the 463// Exec runs one SSH exec command line for an authenticated key. It is the
462// single dispatch path shared by the embedded listener and the system-sshd 464// single dispatch path shared by the embedded listener and the system-sshd
463// forced command (gitbayd shell). Closing revoked kills a git transport. 465// forced command (gitbayd shell). Closing revoked kills a git transport.
464func Exec(cfg config.Config, st *store.Store, packs *packlimit.Limiter, user store.User, key store.SSHKey, term control.Term, cmdline string, 466// packs bounds clones and fetches, and repo download; pushes bounds
467// receive-pack. A nil limiter is no limit.
468func Exec(cfg config.Config, st *store.Store, packs, pushes *packlimit.Limiter, user store.User, key store.SSHKey, term control.Term, cmdline string,
465 stdin io.Reader, stdout, stderr io.Writer, done, stopping, revoked <-chan struct{}) int { 469 stdin io.Reader, stdout, stderr io.Writer, done, stopping, revoked <-chan struct{}) int {
466 if user.Disabled { 470 if user.Disabled {
467 fmt.Fprintln(stderr, "this account is disabled; contact the instance admin") 471 fmt.Fprintln(stderr, "this account is disabled; contact the instance admin")
@@ -479,7 +483,7 @@ func Exec(cfg config.Config, st *store.Store, packs *packlimit.Limiter, user sto
479 if user.Pending { 483 if user.Pending {
480 fmt.Fprintln(stderr, "your account is not active yet: verify your email first") 484 fmt.Fprintln(stderr, "your account is not active yet: verify your email first")
481 } else { 485 } else {
482 code = runGit(cfg, st, packs, user, key.Scope, argv, stdin, stdout, stderr, done, stopping, revoked) 486 code = runGit(cfg, st, packs, pushes, user, key, argv, stdin, stdout, stderr, done, stopping, revoked)
483 } 487 }
484 // A refused push is a refused write, audited like one. runGit 488 // A refused push is a refused write, audited like one. runGit
485 // refuses only with the path as the one argument, so argv[1:] 489 // refuses only with the path as the one argument, so argv[1:]
@@ -512,19 +516,21 @@ func Exec(cfg config.Config, st *store.Store, packs *packlimit.Limiter, user sto
512 Done: done, 516 Done: done,
513 Stopping: stopping, 517 Stopping: stopping,
514 Expires: key.ExpiresAt, 518 Expires: key.ExpiresAt,
519 Packs: packs,
515 } 520 }
516 return control.Dispatch(ctx, argv) 521 return control.Dispatch(ctx, argv)
517} 522}
518 523
519// runGit streams a git transport service after access checks. 524// runGit streams a git transport service after access checks.
520func runGit(cfg config.Config, st *store.Store, packs *packlimit.Limiter, user store.User, scope string, argv []string, 525func runGit(cfg config.Config, st *store.Store, packs, pushes *packlimit.Limiter, user store.User, key store.SSHKey, argv []string,
521 stdin io.Reader, stdout, stderr io.Writer, done, stopping, revoked <-chan struct{}) int { 526 stdin io.Reader, stdout, stderr io.Writer, done, stopping, revoked <-chan struct{}) int {
522 service := argv[0] 527 service, scope := argv[0], key.Scope
523 if len(argv) != 2 { 528 if len(argv) != 2 {
524 fmt.Fprintf(stderr, "usage: %s <path>\n", service) 529 fmt.Fprintf(stderr, "usage: %s <path>\n", service)
525 return protocol.ExitUsage 530 return protocol.ExitUsage
526 } 531 }
527 write := service == "git-receive-pack" 532 write := service == "git-receive-pack"
533 cancel, keepAlive := revoked, 0
528 534
529 repo, err := st.RepoByPath(argv[1]) 535 repo, err := st.RepoByPath(argv[1])
530 if err != nil { 536 if err != nil {
@@ -594,6 +600,25 @@ func runGit(cfg config.Config, st *store.Store, packs *packlimit.Limiter, user s
594 } 600 }
595 } 601 }
596 if write { 602 if write {
603 // Pushes have their own budget, so a clone storm cannot starve
604 // them or the reverse. The slot covers receive-pack and both
605 // hooks: git waits for post-receive (RefsUpdated) before it
606 // exits, and that work is the push's cost. A deploy key is its
607 // own principal, not the account that registered it.
608 principal := "user:" + strconv.FormatInt(user.ID, 10)
609 if policy.IsDeployScope(scope) {
610 principal = "key:" + strconv.FormatInt(key.ID, 10)
611 }
612 release, code := takeSlot(pushes, principal, done, stderr,
613 "the server is busy: it is at its limit of concurrent pushes; try again in a minute")
614 if code != protocol.ExitOK {
615 return code
616 }
617 slotAt := time.Now()
618 // Deferred before Transport runs, so it fires after receive-pack
619 // has exited, on every path: the client hanging up ends its
620 // stdin and receive-pack with it, and a revoked key kills it.
621 defer release()
597 // hookd answers only a hook that names this receive-pack. 622 // hookd answers only a hook that names this receive-pack.
598 token, err := st.CreatePushToken(repo.ID, user.ID, scope) 623 token, err := st.CreatePushToken(repo.ID, user.ID, scope)
599 if err != nil { 624 if err != nil {
@@ -602,25 +627,68 @@ func runGit(cfg config.Config, st *store.Store, packs *packlimit.Limiter, user s
602 } 627 }
603 defer st.DeletePushToken(token) 628 defer st.DeletePushToken(token)
604 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 }
605 } 685 }
606 cancel := revoked
607 if !write { 686 if !write {
608 // Pack generation shares one budget with smart HTTP and git://. 687 // Pack generation shares one budget with smart HTTP and git://.
609 // receive-pack stays outside it: its post-receive runs after the 688 release, code := takeSlot(packs, "user:"+strconv.FormatInt(user.ID, 10), done, stderr,
610 // client has its report, and must not be queued or killed. 689 "the server is busy: it is at its limit of concurrent clones and fetches; try again in a minute")
611 principal := "user:" + strconv.FormatInt(user.ID, 10) 690 if code != protocol.ExitOK {
612 release, err := packs.Acquire(done, principal) 691 return code
613 if err != nil {
614 packs.Refused("ssh", principal, err)
615 }
616 if errors.Is(err, packlimit.ErrBusy) {
617 fmt.Fprintln(stderr, "the server is busy: it is at its limit of concurrent clones and fetches; try again in a minute")
618 return protocol.ExitFailure
619 }
620 if err != nil {
621 // ErrGone: the client left, or the server is restarting.
622 fmt.Fprintln(stderr, "the server is restarting; try again in a minute")
623 return protocol.ExitFailure
624 } 692 }
625 // Deferred before Transport runs, so it fires after git has 693 // Deferred before Transport runs, so it fires after git has
626 // exited and been waited for. 694 // exited and been waited for.
@@ -668,8 +736,58 @@ func runGit(cfg config.Config, st *store.Store, packs *packlimit.Limiter, user s
668 }() 736 }()
669 cancel = kill 737 cancel = kill
670 } 738 }
671 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 {
672 return protocol.ExitFailure 740 return protocol.ExitFailure
673 } 741 }
674 return protocol.ExitOK 742 return protocol.ExitOK
675} 743}
744
745// takeSlot takes a slot from l for principal, waiting until done closes
746// at most. On a refusal it prints busy, or that the server is going
747// away, and returns a nonzero exit.
748func takeSlot(l *packlimit.Limiter, principal string, done <-chan struct{}, stderr io.Writer, busy string) (release func(), code int) {
749 release, err := l.Acquire(done, principal)
750 if err == nil {
751 return release, protocol.ExitOK
752 }
753 l.Refused("ssh", principal, err)
754 if errors.Is(err, packlimit.ErrBusy) {
755 fmt.Fprintln(stderr, busy)
756 } else {
757 // ErrGone: the client left, or the server is restarting.
758 fmt.Fprintln(stderr, "the server is restarting; try again in a minute")
759 }
760 return nil, protocol.ExitFailure
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.
766func 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.
774type 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
781func (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}
internal/sshd/sshd_test.go +2 −2
@@ -64,7 +64,7 @@ func newTestServer(t *testing.T) testServer {
64 64
65 cfg := config.Default() 65 cfg := config.Default()
66 cfg.Server.Root = root 66 cfg.Server.Root = root
67 srv, err := New(cfg, st, nil) 67 srv, err := New(cfg, st, nil, nil)
68 if err != nil { 68 if err != nil {
69 t.Fatal(err) 69 t.Fatal(err)
70 } 70 }
@@ -245,7 +245,7 @@ func TestUnregisteredKeyMessageNamesFingerprintAndHost(t *testing.T) {
245 // The settings link keeps the site URL's scheme and port. 245 // The settings link keeps the site URL's scheme and port.
246 cfg.Server.SiteURL = "http://forge.test:8080/" 246 cfg.Server.SiteURL = "http://forge.test:8080/"
247 cfg.Registration.Mode = "open" 247 cfg.Registration.Mode = "open"
248 srv, err := New(cfg, st, nil) 248 srv, err := New(cfg, st, nil, nil)
249 if err != nil { 249 if err != nil {
250 t.Fatal(err) 250 t.Fatal(err)
251 } 251 }