limits: bound concurrent pushes; repo download under the pack limit !542
29 files changed, +1235 −69
Layout: unified · split
.gitbay/wiki/Admin.org +40 −3
| @@ -357,8 +357,8 @@ push=. | ||
| 357 | 357 | repositories may take; a push may be no larger than what is left. |
| 358 | 358 | - =pack_concurrency= (3), =pack_per_principal= (2), =pack_queue= (32), |
| 359 | 359 | =pack_queue_wait= (="60s"=) — git pack generation (clones, fetches, |
| 360 | =git archive --remote=, web archive downloads) over SSH, smart HTTP | |
| 361 | and git:// shares one | |
| 360 | =git archive --remote=, web archive downloads, =repo download= over | |
| 361 | SSH and the API) over SSH, smart HTTP and git:// shares one | |
| 362 | 362 | budget: this many at once, this many per account (per client |
| 363 | 363 | address when anonymous: an IPv4 address, or an IPv6 /64), and this |
| 364 | 364 | many waiting for at most the wait. Anonymous clients together hold |
| @@ -375,11 +375,48 @@ push=. | ||
| 375 | 375 | for two minutes: a client reading below about 550 B/s, or an HTTP |
| 376 | 376 | request body that takes over two minutes with nothing written back, |
| 377 | 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 | 381 | three counts 0 means the default and a negative value turns that |
| 380 | 382 | bound off. The defaults suit a four-core host; see [[Performance]]. |
| 381 | 383 | With =ssh.mode = "system"= each SSH session is its own process and |
| 382 | 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 | 420 | - =max_pack_bytes= (2 GiB) — the largest pack one push may send, |
| 384 | 421 | enforced as =receive.maxInputSize= and lowered to what an owner's |
| 385 | 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 | 97 | | Control | Status | Evidence | |
| 98 | 98 | |---------------------------------------------+----------+------------------------------------------------------------------| |
| 99 | 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 | 102 | | Service hardening | in place | systemd sandboxing ([[file:03-Deployment.org][3]]) | |
| 102 | 103 | | Backups offsite and append-only | in place | restic with append-only credentials (documented) | |
| 103 | 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 | 16 | | Area | Gap | Severity | |
| 17 | 17 | |-------+-------------------------------------------------------------------------------------------------------------+----------| |
| 18 | 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 | | |
| 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 | | |
| 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 | | |
| 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 | 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 | | |
| 23 | 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 | 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: | ||
| 47 | 47 | full clones of large repositories are CPU-bound in git itself (the 17s |
| 48 | 48 | clone ran git at ~156% CPU). =limits.pack_concurrency= bounds how many |
| 49 | 49 | run 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 | |
| 51 | indexes the pack it is sent, about a core for a large one, and then runs | |
| 52 | its hooks. =limits.push_concurrency= (2) bounds those on a separate | |
| 53 | budget, one per account or deploy key by default, so a clone storm and | |
| 54 | a push storm each leave the other its slots. | |
| 51 | 55 | |
| 52 | 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 | 367 | - *Timing and traffic analysis.* Token comparison is a hash index lookup |
| 368 | 368 | by design, but nothing has been measured. |
| 369 | 369 | - *Denial of service by resource exhaustion* beyond rate. Concurrent |
| 370 | clones, fetches and web archives are bounded by the pack limit | |
| 371 | (#262); pushes are not, beyond =max_pack_bytes= on each one. Nor are | |
| 372 | pathological diffs, deep histories, or zip bombs in LFS. | |
| 370 | clones, fetches, web archives and =repo download= are bounded by the | |
| 371 | pack limit (#262, #308), and pushes by a separate push limit (#308) | |
| 372 | on top of =max_pack_bytes= on each one. Pathological diffs, deep | |
| 373 | histories and zip bombs in LFS are not. | |
| 373 | 374 | |
| 374 | 375 | A sweep is a point in time. This section says what a reader should not |
| 375 | 376 | assume has been checked. |
CHANGELOG.org +22
| @@ -4,6 +4,28 @@ Versioning follows semver from v0.1.0. Database migrations run | ||
| 4 | 4 | automatically on daemon start; upgrade notes appear per release when |
| 5 | 5 | anything 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 | 29 | * v1.40.0 — 2026-09-29 |
| 8 | 30 | |
| 9 | 31 | DKIM 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 | 262 | if packMax > 1 { |
| 263 | 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 | 274 | errCh := make(chan error, 3) |
| 267 | 275 | var sshSrv *sshd.Server |
| 268 | 276 | var sshLn, gitLn net.Listener |
| 269 | 277 | if cfg.SSH.Mode == "embedded" { |
| 270 | srv, err := sshd.New(cfg, st, packs) | |
| 278 | srv, err := sshd.New(cfg, st, packs, pushes) | |
| 271 | 279 | if err != nil { |
| 272 | 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. | |
| 659 | func 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 | 100 | os.Exit(protocol.ExitUsage) |
| 101 | 101 | } |
| 102 | 102 | // Each forced command is its own process, so there is no |
| 103 | // shared pack 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) | |
| 103 | // shared pack or push budget in system mode. | |
| 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 | 105 | st.Close() |
| 106 | 106 | os.Exit(code) |
| 107 | 107 | return nil |
internal/config/config.go +72 −9
| @@ -35,6 +35,21 @@ const ( | ||
| 35 | 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. | |
| 41 | const ( | |
| 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 | ||
| 38 | 53 | type Config struct { |
| 39 | 54 | Server Server `toml:"server"` |
| 40 | 55 | SSH SSH `toml:"ssh"` |
| @@ -235,12 +250,53 @@ type Limits struct { | ||
| 235 | 250 | PackPerPrincipal int `toml:"pack_per_principal"` |
| 236 | 251 | PackQueue int `toml:"pack_queue"` |
| 237 | 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 | 272 | // PackLimits resolves the pack_* settings for packlimit.New. A zero |
| 241 | 273 | // max or per is no bound; an unbounded queue is math.MaxInt, since |
| 242 | 274 | // packlimit reads a zero queue as no queue at all. |
| 243 | 275 | func (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. | |
| 282 | func (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. | |
| 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 | ||
| 299 | func resolveLimits(maxV, perV, queueV int, waitV string, maxDef, perDef, queueDef int, waitDef time.Duration) (max, per, queue int, wait time.Duration) { | |
| 244 | 300 | pick := func(v, def int) int { |
| 245 | 301 | switch { |
| 246 | 302 | case v == 0: |
| @@ -250,16 +306,15 @@ func (l Limits) PackLimits() (max, per, queue int, wait time.Duration) { | ||
| 250 | 306 | } |
| 251 | 307 | return v |
| 252 | 308 | } |
| 253 | queue = pick(l.PackQueue, DefaultPackQueue) | |
| 254 | if l.PackQueue < 0 { | |
| 309 | queue = pick(queueV, queueDef) | |
| 310 | if queueV < 0 { | |
| 255 | 311 | queue = math.MaxInt |
| 256 | 312 | } |
| 257 | wait = DefaultPackQueueWait | |
| 258 | if d, err := time.ParseDuration(l.PackQueueWait); err == nil && d > 0 { | |
| 313 | wait = waitDef | |
| 314 | if d, err := time.ParseDuration(waitV); err == nil && d > 0 { | |
| 259 | 315 | wait = d |
| 260 | 316 | } |
| 261 | return pick(l.PackConcurrency, DefaultPackConcurrency), | |
| 262 | pick(l.PackPerPrincipal, DefaultPackPerPrincipal), queue, wait | |
| 317 | return pick(maxV, maxDef), pick(perV, perDef), queue, wait | |
| 263 | 318 | } |
| 264 | 319 | |
| 265 | 320 | type Mail struct { |
| @@ -619,9 +674,17 @@ func (c Config) Validate() error { | ||
| 619 | 674 | if c.Limits.MaxReposPerUser < 0 || c.Limits.MaxBytesPerUser < 0 || c.Limits.MaxSnippetsPerUser < 0 { |
| 620 | 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 != "" { | |
| 623 | if d, err := time.ParseDuration(w); err != nil || d <= 0 { | |
| 624 | errs = append(errs, fmt.Errorf("limits.pack_queue_wait %q must be a positive duration such as 60s", w)) | |
| 677 | for _, w := range []struct{ name, val string }{ | |
| 678 | {"pack_queue_wait", c.Limits.PackQueueWait}, | |
| 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 | 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. | |
| 67 | func 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 | ||
| 66 | 99 | func TestContradictions(t *testing.T) { |
| 67 | 100 | cases := []struct { |
| 68 | 101 | name string |
| @@ -74,6 +107,21 @@ func TestContradictions(t *testing.T) { | ||
| 74 | 107 | minimal + "\n[limits]\npack_queue_wait = \"soon\"\n", |
| 75 | 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 | 126 | "registration open without smtp", |
| 79 | 127 | minimal + "\n[registration]\nmode = \"open\"\n", |
internal/control/control.go +7
| @@ -14,6 +14,7 @@ import ( | ||
| 14 | 14 | "time" |
| 15 | 15 | |
| 16 | 16 | "gitbay.org/gitbay/internal/config" |
| 17 | "gitbay.org/gitbay/internal/packlimit" | |
| 17 | 18 | "gitbay.org/gitbay/internal/protocol" |
| 18 | 19 | "gitbay.org/gitbay/internal/store" |
| 19 | 20 | ) |
| @@ -67,6 +68,12 @@ type Ctx struct { | ||
| 67 | 68 | // restarting. It closes Done too; a command that ends on Done checks |
| 68 | 69 | // it to say why. |
| 69 | 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 | 79 | // SourceWeb is Ctx.Source for a request from a browser session. Its |
internal/control/download_test.go added +71
| @@ -0,0 +1,71 @@ | ||
| 1 | package control | |
| 2 | ||
| 3 | import ( | |
| 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. | |
| 16 | func 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. | |
| 51 | func 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 @@ | ||
| 1 | 1 | package control |
| 2 | 2 | |
| 3 | 3 | import ( |
| 4 | "errors" | |
| 4 | 5 | "io" |
| 6 | "strconv" | |
| 5 | 7 | "strings" |
| 6 | 8 | |
| 7 | 9 | "gitbay.org/gitbay/internal/gitutil" |
| 10 | "gitbay.org/gitbay/internal/packlimit" | |
| 8 | 11 | "gitbay.org/gitbay/internal/policy" |
| 9 | 12 | "gitbay.org/gitbay/internal/protocol" |
| 10 | 13 | ) |
| @@ -117,6 +120,23 @@ func runRepoDownload(c *Ctx, args []string) int { | ||
| 117 | 120 | // The prefix git puts on every path inside the archive, so unpacking |
| 118 | 121 | // lands in a named directory rather than the current one. |
| 119 | 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 | 140 | if err := gitutil.Archive(dir, ref, prefix, c.Stdout); err != nil { |
| 121 | 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 | 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 | 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 | 35 | // Transport streams one git transport service (upload-pack, receive-pack, |
| 36 | 36 | // upload-archive). extraEnv entries are appended to the process environment; |
| 37 | 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 | 40 | // and everything it started; a push killed before its pre-receive hook |
| 40 | 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 | 43 | var args []string |
| 43 | 44 | switch service { |
| 44 | 45 | case "git-upload-pack", "git-receive-pack", "git-upload-archive": |
| 45 | 46 | if service == "git-receive-pack" && maxPack > 0 { |
| 46 | 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 | 54 | if service == "git-upload-pack" { |
| 49 | 55 | // Keepalives while pack-objects is still counting keep a |
| 50 | 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 | 19 | in, w := io.Pipe() |
| 20 | 20 | cancel := make(chan struct{}) |
| 21 | 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 | 23 | close(cancel) |
| 24 | 24 | time.AfterFunc(500*time.Millisecond, func() { w.Close() }) |
| 25 | 25 | select { |
internal/hookd/hookd.go +35
| @@ -17,6 +17,7 @@ import ( | ||
| 17 | 17 | "os" |
| 18 | 18 | "path/filepath" |
| 19 | 19 | "strings" |
| 20 | "sync" | |
| 20 | 21 | |
| 21 | 22 | "gitbay.org/gitbay/internal/config" |
| 22 | 23 | "gitbay.org/gitbay/internal/control" |
| @@ -139,6 +140,7 @@ func (s *Server) handle(conn net.Conn) { | ||
| 139 | 140 | } |
| 140 | 141 | switch req.Hook { |
| 141 | 142 | case "pre-receive": |
| 143 | preReceiveStarted(req.Token) | |
| 142 | 144 | s.preReceive(req, dec, enc) |
| 143 | 145 | case "post-receive": |
| 144 | 146 | s.postReceive(req) |
| @@ -324,3 +326,36 @@ func WriteHookScripts(hooksDir, gitbaydPath string) error { | ||
| 324 | 326 | } |
| 325 | 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 | 9 | "path/filepath" |
| 10 | 10 | "strings" |
| 11 | 11 | "testing" |
| 12 | "time" | |
| 12 | 13 | |
| 13 | 14 | "gitbay.org/gitbay/internal/config" |
| 14 | 15 | "gitbay.org/gitbay/internal/policy" |
| @@ -198,3 +199,33 @@ func TestRefusedPushCapsRefs(t *testing.T) { | ||
| 198 | 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/httpd/api.go +9
| @@ -77,6 +77,7 @@ func (s *Server) apiCmd(w http.ResponseWriter, r *http.Request) { | ||
| 77 | 77 | Expires: tok.ExpiresAt, |
| 78 | 78 | Done: s.until(r), |
| 79 | 79 | Stopping: s.stopping, |
| 80 | Packs: s.packs, | |
| 80 | 81 | } |
| 81 | 82 | code := control.Dispatch(ctx, req.Argv) |
| 82 | 83 | |
| @@ -96,12 +97,20 @@ func (s *Server) apiCmd(w http.ResponseWriter, r *http.Request) { | ||
| 96 | 97 | body["stderr"] = msg |
| 97 | 98 | } |
| 98 | 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 | 104 | w.WriteHeader(status) |
| 100 | 105 | json.NewEncoder(w).Encode(body) |
| 101 | 106 | } |
| 102 | 107 | |
| 103 | 108 | // statusForExit maps a command's exit code onto an HTTP status, shared by |
| 104 | 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. | |
| 112 | const busyRetryAfter = "60" | |
| 113 | ||
| 105 | 114 | func statusForExit(code int) int { |
| 106 | 115 | switch code { |
| 107 | 116 | case protocol.ExitOK: |
internal/httpd/apibusy_test.go added +52
| @@ -0,0 +1,52 @@ | ||
| 1 | package httpd | |
| 2 | ||
| 3 | import ( | |
| 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. | |
| 16 | func 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 | 64 | ReadOnly: true, |
| 65 | 65 | Done: s.until(r), |
| 66 | 66 | Stopping: s.stopping, |
| 67 | Packs: s.packs, | |
| 67 | 68 | } |
| 68 | 69 | code := control.Dispatch(ctx, argv) |
| 69 | 70 | |
| @@ -81,6 +82,15 @@ func (s *Server) apiRead(w http.ResponseWriter, r *http.Request) { | ||
| 81 | 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 | 94 | // Responses are authorized per account, so the ETag is salted with the |
| 85 | 95 | // caller: two users asking the same question may get different answers, |
| 86 | 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 | |
| 2 | // and upload-archive over SSH, smart HTTP and git:// draw on one | |
| 3 | // budget: a global cap, a cap per principal (an account, or a client | |
| 4 | // address on the anonymous transports), and a bounded queue whose | |
| 5 | // waiters give up after a fixed wait or when the client goes away. | |
| 1 | // Package packlimit bounds concurrent git processes. upload-pack and | |
| 2 | // upload-archive over SSH, smart HTTP and git:// draw on one budget, | |
| 3 | // receive-pack on a second: each a global cap, a cap per principal (an | |
| 4 | // account, a deploy key, or a client address on the anonymous | |
| 5 | // transports), and a bounded queue whose waiters give up after a fixed | |
| 6 | // wait or when the client goes away. | |
| 6 | 7 | // Waiters are not served in order; a new arrival can take a freed slot |
| 7 | 8 | // ahead of them, and the wait bounds how long any one of them waits. |
| 8 | 9 | package packlimit |
| @@ -23,7 +24,9 @@ var ( | ||
| 23 | 24 | |
| 24 | 25 | type Limiter struct { |
| 25 | 26 | max, per, queue int |
| 27 | perQueue int // waiting, per principal; per unless set | |
| 26 | 28 | wait time.Duration |
| 29 | name string // what is limited, for the refusal log | |
| 27 | 30 | |
| 28 | 31 | // Principals starting with class may hold at most classCap slots |
| 29 | 32 | // between them; classCap 0 is no class cap. |
| @@ -45,7 +48,7 @@ func New(max, per, queue int, wait time.Duration) *Limiter { | ||
| 45 | 48 | if max <= 0 { |
| 46 | 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 | 52 | held: map[string]int{}, waiting: map[string]int{}, changed: make(chan struct{}), |
| 50 | 53 | warned: map[string]time.Time{}} |
| 51 | 54 | } |
| @@ -71,10 +74,29 @@ func (l *Limiter) Refused(transport, principal string, err error) { | ||
| 71 | 74 | if errors.Is(err, ErrGone) { |
| 72 | 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 | 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. | |
| 83 | func (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. | |
| 93 | func (l *Limiter) CapQueue(n int) { | |
| 94 | if l == nil { | |
| 95 | return | |
| 96 | } | |
| 97 | l.perQueue = n | |
| 98 | } | |
| 99 | ||
| 78 | 100 | // CapClass caps the slots that principals starting with prefix may hold |
| 79 | 101 | // between them. Call it before the limiter is in use. |
| 80 | 102 | func (l *Limiter) CapClass(prefix string, n int) { |
| @@ -116,7 +138,7 @@ func (l *Limiter) Acquire(done <-chan struct{}, principal string) (release func( | ||
| 116 | 138 | l.mu.Unlock() |
| 117 | 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 | 142 | l.mu.Unlock() |
| 121 | 143 | return nil, ErrBusy |
| 122 | 144 | } |
internal/packlimit/packlimit_test.go +52
| @@ -313,3 +313,55 @@ func TestRefusedLogsOncePerTransport(t *testing.T) { | ||
| 313 | 313 | var none *Limiter |
| 314 | 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. | |
| 319 | func 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. | |
| 343 | func 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 | 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 | 3 | import ( |
| 4 | 4 | "io" |
| 5 | "strings" | |
| 5 | 6 | "testing" |
| 6 | 7 | "time" |
| 7 | 8 | ) |
| @@ -38,3 +39,37 @@ func TestWatchWithoutLimit(t *testing.T) { | ||
| 38 | 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 +181 −7
| @@ -2,9 +2,11 @@ package sshd | ||
| 2 | 2 | |
| 3 | 3 | import ( |
| 4 | 4 | "bytes" |
| 5 | "errors" | |
| 5 | 6 | "io" |
| 6 | 7 | "os" |
| 7 | 8 | "path/filepath" |
| 9 | "strconv" | |
| 8 | 10 | "strings" |
| 9 | 11 | "testing" |
| 10 | 12 | "time" |
| @@ -56,7 +58,7 @@ func TestRefusedPushIsAudited(t *testing.T) { | ||
| 56 | 58 | bob.Pending = pending |
| 57 | 59 | key := store.SSHKey{Scope: "full", Fingerprint: "SHA256:test"} |
| 58 | 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 | 62 | strings.NewReader(""), &out, &errOut, nil, nil, nil) |
| 61 | 63 | if code != protocol.ExitDenied { |
| 62 | 64 | t.Fatalf("pending %v: exit %d: %s", pending, code, errOut.String()) |
| @@ -83,7 +85,7 @@ func TestCloneRefusedWhenPackSlotsAreFull(t *testing.T) { | ||
| 83 | 85 | key := store.SSHKey{Scope: "full", Fingerprint: "SHA256:test"} |
| 84 | 86 | for _, service := range []string{"git-upload-pack", "git-upload-archive"} { |
| 85 | 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 | 89 | strings.NewReader(""), &out, &errOut, nil, nil, nil) |
| 88 | 90 | if code != protocol.ExitFailure || !strings.Contains(errOut.String(), "busy") { |
| 89 | 91 | t.Fatalf("%s: exit %d: %q", service, code, errOut.String()) |
| @@ -106,7 +108,7 @@ func TestPushBypassesPackLimitAndCloneReleasesSlot(t *testing.T) { | ||
| 106 | 108 | packs := packlimit.New(1, 0, 0, time.Second) |
| 107 | 109 | |
| 108 | 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 | 112 | strings.NewReader("0000"), &out, &errOut, nil, nil, nil); code != protocol.ExitOK { |
| 111 | 113 | t.Fatalf("clone: exit %d: %s", code, errOut.String()) |
| 112 | 114 | } |
| @@ -118,7 +120,7 @@ func TestPushBypassesPackLimitAndCloneReleasesSlot(t *testing.T) { | ||
| 118 | 120 | |
| 119 | 121 | out.Reset() |
| 120 | 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 | 124 | strings.NewReader("0000"), &out, &errOut, nil, nil, nil); code != protocol.ExitOK { |
| 123 | 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 | 161 | packs := packlimit.New(1, 0, 0, time.Second) |
| 160 | 162 | codec := make(chan int, 1) |
| 161 | 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 | 165 | silentStdin(t), stdout, io.Discard, done, stopping, revoked) |
| 164 | 166 | }() |
| 165 | 167 | select { |
| @@ -208,7 +210,7 @@ func TestCloneRunsOnDuringRestart(t *testing.T) { | ||
| 208 | 210 | key := store.SSHKey{Scope: "full", Fingerprint: "SHA256:test"} |
| 209 | 211 | packs := packlimit.New(1, 0, 0, time.Second) |
| 210 | 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 | 214 | strings.NewReader("0000"), io.Discard, &errOut, closed(), closed(), nil); code != protocol.ExitOK { |
| 213 | 215 | t.Fatalf("exit %d: %s", code, errOut.String()) |
| 214 | 216 | } |
| @@ -233,9 +235,181 @@ func TestRefusedCloneStaysOffLimiter(t *testing.T) { | ||
| 233 | 235 | defer hold() |
| 234 | 236 | key := store.SSHKey{Scope: "full", Fingerprint: "SHA256:test"} |
| 235 | 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 | 239 | strings.NewReader(""), &out, &errOut, nil, nil, nil) |
| 238 | 240 | if code != protocol.ExitNotFound || strings.Contains(errOut.String(), "busy") { |
| 239 | 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. | |
| 247 | func 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. | |
| 259 | func 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. | |
| 274 | func 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. | |
| 350 | func 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. | |
| 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() | |
| 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. | |
| 400 | func 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 @@ | ||
| 3 | 3 | package sshd |
| 4 | 4 | |
| 5 | 5 | import ( |
| 6 | "bytes" | |
| 6 | 7 | "context" |
| 7 | 8 | "crypto/ed25519" |
| 8 | 9 | "crypto/rand" |
| @@ -39,6 +40,7 @@ type Server struct { | ||
| 39 | 40 | cfg config.Config |
| 40 | 41 | st *store.Store |
| 41 | 42 | packs *packlimit.Limiter |
| 43 | pushes *packlimit.Limiter | |
| 42 | 44 | sshCfg *ssh.ServerConfig |
| 43 | 45 | authLimiter *rateLimiter |
| 44 | 46 | sessions sync.WaitGroup // accepted connections still being served |
| @@ -70,8 +72,8 @@ func (c *conn) cut() { | ||
| 70 | 72 | c.net.Close() |
| 71 | 73 | } |
| 72 | 74 | |
| 73 | func New(cfg config.Config, st *store.Store, packs *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{})} | |
| 75 | func New(cfg config.Config, st *store.Store, packs, pushes *packlimit.Limiter) (*Server, error) { | |
| 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 | 78 | sc := &ssh.ServerConfig{ |
| 77 | 79 | PublicKeyCallback: s.authenticate, |
| @@ -425,7 +427,7 @@ func (s *Server) runExec(c *conn, sconn *ssh.ServerConn, ch ssh.Channel, term co | ||
| 425 | 427 | return protocol.ExitDenied |
| 426 | 428 | } |
| 427 | 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 | 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 | 463 | // Exec runs one SSH exec command line for an authenticated key. It is the |
| 462 | 464 | // single dispatch path shared by the embedded listener and the system-sshd |
| 463 | 465 | // forced command (gitbayd shell). Closing revoked kills a git transport. |
| 464 | func 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. | |
| 468 | func Exec(cfg config.Config, st *store.Store, packs, pushes *packlimit.Limiter, user store.User, key store.SSHKey, term control.Term, cmdline string, | |
| 465 | 469 | stdin io.Reader, stdout, stderr io.Writer, done, stopping, revoked <-chan struct{}) int { |
| 466 | 470 | if user.Disabled { |
| 467 | 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 | 483 | if user.Pending { |
| 480 | 484 | fmt.Fprintln(stderr, "your account is not active yet: verify your email first") |
| 481 | 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 | 488 | // A refused push is a refused write, audited like one. runGit |
| 485 | 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 | 516 | Done: done, |
| 513 | 517 | Stopping: stopping, |
| 514 | 518 | Expires: key.ExpiresAt, |
| 519 | Packs: packs, | |
| 515 | 520 | } |
| 516 | 521 | return control.Dispatch(ctx, argv) |
| 517 | 522 | } |
| 518 | 523 | |
| 519 | 524 | // runGit streams a git transport service after access checks. |
| 520 | func runGit(cfg config.Config, st *store.Store, packs *packlimit.Limiter, user store.User, scope string, argv []string, | |
| 525 | func runGit(cfg config.Config, st *store.Store, packs, pushes *packlimit.Limiter, user store.User, key store.SSHKey, argv []string, | |
| 521 | 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 | 528 | if len(argv) != 2 { |
| 524 | 529 | fmt.Fprintf(stderr, "usage: %s <path>\n", service) |
| 525 | 530 | return protocol.ExitUsage |
| 526 | 531 | } |
| 527 | 532 | write := service == "git-receive-pack" |
| 533 | cancel, keepAlive := revoked, 0 | |
| 528 | 534 | |
| 529 | 535 | repo, err := st.RepoByPath(argv[1]) |
| 530 | 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 | 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 | 622 | // hookd answers only a hook that names this receive-pack. |
| 598 | 623 | token, err := st.CreatePushToken(repo.ID, user.ID, scope) |
| 599 | 624 | if err != nil { |
| @@ -602,25 +627,68 @@ func runGit(cfg config.Config, st *store.Store, packs *packlimit.Limiter, user s | ||
| 602 | 627 | } |
| 603 | 628 | defer st.DeletePushToken(token) |
| 604 | 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 | 686 | if !write { |
| 608 | 687 | // Pack generation shares one budget with smart HTTP and git://. |
| 609 | // receive-pack stays outside it: its post-receive runs after the | |
| 610 | // client has its report, and must not be queued or killed. | |
| 611 | principal := "user:" + strconv.FormatInt(user.ID, 10) | |
| 612 | release, err := packs.Acquire(done, principal) | |
| 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 | |
| 688 | release, code := takeSlot(packs, "user:"+strconv.FormatInt(user.ID, 10), done, stderr, | |
| 689 | "the server is busy: it is at its limit of concurrent clones and fetches; try again in a minute") | |
| 690 | if code != protocol.ExitOK { | |
| 691 | return code | |
| 624 | 692 | } |
| 625 | 693 | // Deferred before Transport runs, so it fires after git has |
| 626 | 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 | 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 | 740 | return protocol.ExitFailure |
| 673 | 741 | } |
| 674 | 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. | |
| 748 | func 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. | |
| 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 | } | |
internal/sshd/sshd_test.go +2 −2
| @@ -64,7 +64,7 @@ func newTestServer(t *testing.T) testServer { | ||
| 64 | 64 | |
| 65 | 65 | cfg := config.Default() |
| 66 | 66 | cfg.Server.Root = root |
| 67 | srv, err := New(cfg, st, nil) | |
| 67 | srv, err := New(cfg, st, nil, nil) | |
| 68 | 68 | if err != nil { |
| 69 | 69 | t.Fatal(err) |
| 70 | 70 | } |
| @@ -245,7 +245,7 @@ func TestUnregisteredKeyMessageNamesFingerprintAndHost(t *testing.T) { | ||
| 245 | 245 | // The settings link keeps the site URL's scheme and port. |
| 246 | 246 | cfg.Server.SiteURL = "http://forge.test:8080/" |
| 247 | 247 | cfg.Registration.Mode = "open" |
| 248 | srv, err := New(cfg, st, nil) | |
| 248 | srv, err := New(cfg, st, nil, nil) | |
| 249 | 249 | if err != nil { |
| 250 | 250 | t.Fatal(err) |
| 251 | 251 | } |