packlimit: one limit on pack generation across transports !506

merged merged by cmc on 2026-09-28 23:13 UTC · krz/gitbay:pack-limit into main

39 files changed, +1831 −96

Layout: unified · split

.gitbay/wiki/Admin.org +42 −5
@@ -208,7 +208,36 @@ push=.
208 Organizations are not capped. 208 Organizations are not capped.
209- =max_bytes_per_user= (0, unlimited) — disk the account's own 209- =max_bytes_per_user= (0, unlimited) — disk the account's own
210 repositories may take; a push may be no larger than what is left. 210 repositories may take; a push may be no larger than what is left.
211- =max_pack_bytes=, =ssh_auth_rate= — reserved, not yet enforced. 211- =pack_concurrency= (3), =pack_per_principal= (2), =pack_queue= (32),
212 =pack_queue_wait= (="60s"=) — git pack generation (clones, fetches,
213 =git archive --remote=, web archive downloads) over SSH, smart HTTP
214 and git:// shares one
215 budget: this many at once, this many per account (per client
216 address when anonymous: an IPv4 address, or an IPv6 /64), and this
217 many waiting for at most the wait. Anonymous clients together hold
218 at most =pack_concurrency= − 1 slots when it is above 1: an account
219 (SSH key, bearer token or web session) can take the last slot when it
220 is free, and anonymous clients cannot hold it.
221 Past that an SSH client gets "the server is busy…" and exit 1, HTTP
222 gets 503 with =Retry-After: 30=, git:// an =ERR= line; the daemon
223 logs a =pack limit= warning naming the transport and whether the
224 client was signed in, at most once a minute per transport. An SSH client
225 that disconnects while queued leaves the queue; an HTTP or git:// one
226 keeps its place until the wait runs out. A running clone is killed
227 when its client disconnects, or when no write to the client completes
228 for two minutes: a client reading below about 550 B/s, or an HTTP
229 request body that takes over two minutes with nothing written back,
230 is cut. Ref listings (info/refs, protocol v2
231 =ls-refs=), pushes and =repo download= are outside the budget. For the
232 three counts 0 means the default and a negative value turns that
233 bound off. The defaults suit a four-core host; see [[Performance]].
234 With =ssh.mode = "system"= each SSH session is its own process and
235 SSH clones are not counted.
236- =max_pack_bytes= (2 GiB) — the largest pack one push may send,
237 enforced as =receive.maxInputSize= and lowered to what an owner's
238 storage quota has left.
239- =ssh_auth_rate= (10) — SSH authentication failures per client address
240 per minute on the embedded listener; see below.
212 241
213** [git_daemon] 242** [git_daemon]
214- =enabled= (false), =port= (9418) — the anonymous =git://= listener. 243- =enabled= (false), =port= (9418) — the anonymous =git://= listener.
@@ -258,8 +287,12 @@ meaningful — an unverified address never produces a =verified= badge.
258The audit log is the security feed (events are the product feed): every 287The audit log is the security feed (events are the product feed): every
259successful mutating command with its argv and source credential (SSH key 288successful mutating command with its argv and source credential (SSH key
260fingerprint or API), every refused one (exit 3 or 4) as =refused 289fingerprint or API), every refused one (exit 3 or 4) as =refused
261<command>=, refused pushes as =refused git-receive-pack=, registrations, 290<command>=, refused pushes as =refused git-receive-pack= (the access
262admin actions, force-pushes, and auth failures/throttling. A refusal row 291check) or =refused push= (a branch or tag rule, a release anchor or an
292unsigned commit, with the repository and ref names), hook socket
293requests failing the peer or push-token check as =refused hook= (never
294the token), registrations, admin actions, force-pushes, and auth
295failures/throttling. A refusal row
263keeps the flag names and the first positional, not the values. 296keeps the flag names and the first positional, not the values.
264Refusals are recorded up to ten a minute per account and 600 a minute 297Refusals are recorded up to ten a minute per account and 600 a minute
265across the instance; past either, one =refused.throttled= row stands 298across the instance; past either, one =refused.throttled= row stands
@@ -272,7 +305,10 @@ Each row carries the SHA-256 of the row before it. =gitbayd admin audit
272verify= opens the store as other admin commands do, applying pending 305verify= opens the store as other admin commands do, applying pending
273migrations, so run it with the binary that matches the daemon. It 306migrations, so run it with the binary that matches the daemon. It
274recomputes the chain and exits 1 naming the first row that was 307recomputes the chain and exits 1 naming the first row that was
275edited or whose predecessor was removed. Retention removing the oldest 308edited or whose predecessor was removed. The chain is unkeyed: whoever
309can write the database can recompute every hash after an edit, and
310verify then finds nothing. It catches an edit only when the later
311hashes were not recomputed. Retention removing the oldest
276rows is not a break. Rows written before the chain existed are counted 312rows is not a break. Rows written before the chain existed are counted
277and skipped; when every row is such a row, verify warns and exits 1, 313and skipped; when every row is such a row, verify warns and exits 1,
278since clearing the hash columns looks the same. After an upgrade that 314since clearing the hash columns looks the same. After an upgrade that
@@ -282,7 +318,8 @@ Removing the newest rows leaves no break, and neither do rows written
282afterwards under the freed ids. The database cannot show either. The 318afterwards under the freed ids. The database cannot show either. The
283daemon logs every row it writes to its journal, outside the database 319daemon logs every row it writes to its journal, outside the database
284(=journalctl -u gitbayd -g 'INFO audit '=), and verify prints the last 320(=journalctl -u gitbayd -g 'INFO audit '=), and verify prints the last
285id and hash: compare them with the newest journal line. Rows written by 321id and hash: comparing them with the newest journal line is the check
322for any change, recomputed hashes included. Rows written by
286host =gitbayd admin= commands, and by =gitbayd shell= when =ssh.mode = 323host =gitbayd admin= commands, and by =gitbayd shell= when =ssh.mode =
287"system"=, are not copied to the journal. 324"system"=, are not copied to the journal.
288 325
.gitbay/wiki/Architecture/09-Controls.org +2 −2
@@ -69,7 +69,7 @@ chapter names of OWASP ASVS 4.0 where one fits.
69| Security-relevant writes audited | in place | every successful mutating command (=control.go=) | 69| Security-relevant writes audited | in place | every successful mutating command (=control.go=) |
70| Authentication failures audited | in place | =auth.failed=, =auth.throttled= | 70| Authentication failures audited | in place | =auth.failed=, =auth.throttled= |
71| Denied attempts audited | in place | refused mutating commands and pushes, ten a minute per actor, 600 in all (=internal/control/auditrefusal.go=) | 71| Denied attempts audited | in place | refused mutating commands and pushes, ten a minute per actor, 600 in all (=internal/control/auditrefusal.go=) |
72| Audit log tamper resistance | partial | hash chain checked by =gitbayd admin audit verify=; every row the daemon writes copied to its journal; the table is writable by the daemon user, and removing the newest rows (or reusing their ids) shows only by comparing verify's last id and hash with the journal | 72| Audit log tamper resistance | partial | unkeyed hash chain checked by =gitbayd admin audit verify=; every row the daemon writes copied to its journal; the table is writable by the daemon user, who can recompute the chain after an edit, so comparing verify's last id and hash with the journal is the check for any change |
73 73
74** Communications and integrations (V9, V10, V12) 74** Communications and integrations (V9, V10, V12)
75 75
@@ -96,7 +96,7 @@ chapter names of OWASP ASVS 4.0 where one fits.
96| Control | Status | Evidence | 96| Control | Status | Evidence |
97|---------------------------------------------+----------+------------------------------------------------------------------| 97|---------------------------------------------+----------+------------------------------------------------------------------|
98| Rate limits on API and writes | in place | [[file:05-Identity-and-Access.org][5. Rate limits]] | 98| Rate limits on API and writes | in place | [[file:05-Identity-and-Access.org][5. Rate limits]] |
99| Concurrency limit on git pack generation | gap | #262 | 99| 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| Service hardening | in place | systemd sandboxing ([[file:03-Deployment.org][3]]) | 100| Service hardening | in place | systemd sandboxing ([[file:03-Deployment.org][3]]) |
101| Backups offsite and append-only | in place | restic with append-only credentials (documented) | 101| Backups offsite and append-only | in place | restic with append-only credentials (documented) |
102| Restore tested | gap | #259 | 102| Restore tested | gap | #259 |
.gitbay/wiki/Architecture/10-Known-Gaps.org +6 −3
@@ -13,15 +13,18 @@ what the 2026-09-27 review found; remove a row when its issue closes.
13| #259 | Recovery | No restore has been exercised; the drill is written (Admin wiki) and not yet run | high | 13| #259 | Recovery | No restore has been exercised; the drill is written (Admin wiki) and not yet run | high |
14| #260 | CI network | Builds share the runner's source address; no egress policy | medium | 14| #260 | CI network | Builds share the runner's source address; no egress policy | medium |
15| #261 | Various | Migration foreign-key check after commit; three web writes bypass dispatch; documentation drift | medium | 15| #261 | Various | Migration foreign-key check after commit; three web writes bypass dispatch; documentation drift | medium |
16| #262 | Availability | No limit on concurrent git pack generation | high |
17| #298 | SSRF | =repo import --from= fetches without an address check | medium |
18| #297 | Credentials | A browser session can mint tokens and keys that outlive it | low | 16| #297 | Credentials | A browser session can mint tokens and keys that outlive it | low |
17| #298 | SSRF | =repo import --from= fetches without an address check | medium |
19 18
20* Not filed 19* Not filed
21 20
22| Area | Gap | Severity | 21| Area | Gap | Severity |
23|-------+-------------------------------------------------------------------------------------------------------------+----------| 22|-------+-------------------------------------------------------------------------------------------------------------+----------|
24| Audit | Removing the newest audit rows, or writing new rows under their freed ids, is not 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 | 23| 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 |
24| 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 |
25| Availability | Pushes have no concurrency limit; =max_pack_bytes= bounds each one, not how many run at once | medium |
26| Availability | =repo download= (SSH, API) runs =git archive= outside the pack limit; only its two-minute deadline and 512 MiB cap bound it | low |
27| 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 |
25 28
26* Questions an auditor will ask that have no answer yet 29* Questions an auditor will ask that have no answer yet
27 30
.gitbay/wiki/Performance.org +17 −2
@@ -45,5 +45,20 @@ this scale never appears in a profile.
45 45
46The practical ceiling on this hardware is concurrent pack generation: 46The 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). A busier instance would scale that with 48clone ran git at ~156% CPU). =limits.pack_concurrency= bounds how many
49cores, not with changes to gitbay. 49run at once across SSH, HTTP and git://, with a queue behind it (see
50[[Admin]], =[limits]=).
51
52* Concurrent clones
53
54The defaults (=pack_concurrency= 3, =pack_per_principal= 2,
55=pack_queue= 32, =pack_queue_wait= 60s) are set for a four-core host
56from the single-clone figure above; no concurrent-clone measurement
57backs them yet. =deploy/clonebench.sh <clone-url> <n>= starts n full
58clones at once from a machine other than the server. Anonymous clones
59together hold at most =pack_concurrency= − 1 slots, and clones from one
60address or account count against one principal, so an anonymous run
61measures those caps rather than the global one. Run it as signed-in
62clients (SSH keys, or HTTP with bearer tokens) from more than one
63account, or against an instance with =pack_per_principal = -1= and
64read the result as the anonymous class cap.
.gitbay/wiki/Threat-Model.org +13 −9
@@ -249,8 +249,10 @@ is recorded here rather than in a closed issue:
249 escape what is there. #144 covers the missing isolation. 249 escape what is there. #144 covers the missing isolation.
250- *Timing and traffic analysis.* Token comparison is a hash index lookup 250- *Timing and traffic analysis.* Token comparison is a hash index lookup
251 by design, but nothing has been measured. 251 by design, but nothing has been measured.
252- *Denial of service by resource exhaustion* beyond rate: large pushes, 252- *Denial of service by resource exhaustion* beyond rate. Concurrent
253 pathological diffs, deep histories, zip bombs in LFS. 253 clones, fetches and web archives are bounded by the pack limit
254 (#262); pushes are not, beyond =max_pack_bytes= on each one. Nor are
255 pathological diffs, deep histories, or zip bombs in LFS.
254 256
255A sweep is a point in time. This section says what a reader should not 257A sweep is a point in time. This section says what a reader should not
256assume has been checked. 258assume has been checked.
@@ -271,13 +273,15 @@ assume has been checked.
271 present as objects no archived ref names, and a repository's refs may 273 present as objects no archived ref names, and a repository's refs may
272 be newer than the database snapshot (see [[Admin]]). 274 be newer than the database snapshot (see [[Admin]]).
273- The audit log lives in the database the daemon writes, so anyone with 275- The audit log lives in the database the daemon writes, so anyone with
274 the daemon user's access can change it. The hash chain makes an edited 276 the daemon user's access can change it. The hash chain is unkeyed:
275 or removed row show as a break under =gitbayd admin audit verify=, 277 whoever can write the database can edit a row and recompute every
276 except at the end: removing the newest rows, and writing new rows 278 later hash. =gitbayd admin audit verify= catches an edited or removed
277 under their freed ids, leaves a valid chain. Only comparing verify's 279 row only when the later hashes were not recomputed, and never catches
278 last id and hash with the daemon's journal copy shows it, and rows 280 removing the newest rows or writing new rows under their freed ids.
279 written outside the daemon (=gitbayd shell= under =ssh.mode = 281 Comparing verify's last id and hash with the daemon's journal copy is
280 "system"=, host =gitbayd admin= commands) have no journal copy. 282 the check for any change; rows written outside the daemon (=gitbayd
283 shell= under =ssh.mode = "system"=, host =gitbayd admin= commands)
284 have no journal copy.
281- A global signature-verification epoch over-invalidates the cache on any 285- A global signature-verification epoch over-invalidates the cache on any
282 trust-input change. Correct, not a leak; a performance tradeoff. 286 trust-input change. Correct, not a leak; a performance tradeoff.
283- A build's secrets are environment variables inside its container, so 287- A build's secrets are environment variables inside its container, so
CHANGELOG.org +32 −2
@@ -99,8 +99,11 @@ missing, =gitbayd admin backup --verify <archive>= names it, and
99 before it (migration 0064), and the daemon logs a copy of every row it 99 before it (migration 0064), and the daemon logs a copy of every row it
100 writes to its journal. =gitbayd admin audit verify= prints the row 100 writes to its journal. =gitbayd admin audit verify= prints the row
101 count and the last id and hash, and exits 1 naming the first row that 101 count and the last id and hash, and exits 1 naming the first row that
102 was edited or whose predecessor was removed. Removing the newest rows 102 was edited or whose predecessor was removed. The chain is unkeyed:
103 shows only by comparing that last id and hash with the journal (#275). 103 someone who can write the database can recompute the later hashes,
104 and verify catches an edit only when they were not recomputed.
105 Comparing verify's last id and hash with the journal is the check for
106 any change, including removing the newest rows (#275).
104- Refused mutating commands (exit 3 or 4) are audited as =refused 107- Refused mutating commands (exit 3 or 4) are audited as =refused
105 <command>=, and refused pushes as =refused git-receive-pack=, keeping 108 <command>=, and refused pushes as =refused git-receive-pack=, keeping
106 flag names and the target but no values; ten a minute per account and 109 flag names and the target but no values; ten a minute per account and
@@ -108,6 +111,10 @@ missing, =gitbayd admin backup --verify <archive>= names it, and
108 =refused.throttled= row stands for the rest of the minute. Under 111 =refused.throttled= row stands for the rest of the minute. Under
109 =ssh.mode = "system"= each =gitbayd shell= connection counts 112 =ssh.mode = "system"= each =gitbayd shell= connection counts
110 separately (#275). 113 separately (#275).
114- Pushes refused in pre-receive (a branch or tag rule, a release
115 anchor, an unsigned commit) are audited as =refused push= with the
116 repository and ref names, and hook socket requests failing the peer
117 or push-token check as =refused hook=, under the same caps (#275).
111- Audit retention deletes by id, up to the newest row older than the 118- Audit retention deletes by id, up to the newest row older than the
112 retention, so a clock step back cannot leave a gap in the chain (#275). 119 retention, so a clock step back cannot leave a gap in the chain (#275).
113- =dashboard= and =feed= print activity as sentences 120- =dashboard= and =feed= print activity as sentences
@@ -217,6 +224,29 @@ missing, =gitbayd admin backup --verify <archive>= names it, and
217 repository on the host. (#259) 224 repository on the host. (#259)
218- =gitbayd admin secrets init= and =rotate= hold an flock on =<key 225- =gitbayd admin secrets init= and =rotate= hold an flock on =<key
219 file>.lock=, so two runs at once serialize. (#273) 226 file>.lock=, so two runs at once serialize. (#273)
227- Git pack generation (clones, fetches, =git archive --remote=) over
228 SSH, smart HTTP and git:// now shares one concurrency budget:
229 =limits.pack_concurrency= (3), =pack_per_principal= (2), =pack_queue=
230 (32) and =pack_queue_wait= (60s). Past the queue an SSH client sees
231 "the server is busy…" and exits 1, HTTP gets 503 with
232 =Retry-After: 30=, and git:// gets an =ERR= line. *Operators:* the
233 defaults are tuned for a four-core host; set the three counts to -1
234 to turn the limit off. A running clone is killed when no write to
235 its client completes for two minutes: a client reading below about
236 550 B/s, or an HTTP request body that takes over two minutes with
237 nothing written back, is cut. Ref listings, pushes and =repo
238 download= are unaffected. Under =ssh.mode = "system"= SSH clones are not counted,
239 since each session is its own process (#262).
240- Anonymous clones are counted per IPv4 address or IPv6 /64, and
241 together hold at most =pack_concurrency= − 1 slots: an account can
242 take the last slot when it is free, and anonymous clients cannot hold
243 it (#262).
244- A request turned away by the pack limit logs a warning naming the
245 transport and whether the client was signed in, never its address,
246 at most once a minute per transport (#262).
247- Web archive downloads (=/{owner}/{repo}/archive/{ref}.tar.gz=) take
248 a pack slot, answer 503 with =Retry-After: 30= when none is free, and
249 are killed when the client leaves or stops reading (#262).
220 250
221* v1.36.0 — 2026-09-23 251* v1.36.0 — 2026-09-23
222 252
cmd/gitbayd/main.go +13 −3
@@ -32,6 +32,7 @@ import (
32 "gitbay.org/gitbay/internal/httpd" 32 "gitbay.org/gitbay/internal/httpd"
33 "gitbay.org/gitbay/internal/mirror" 33 "gitbay.org/gitbay/internal/mirror"
34 "gitbay.org/gitbay/internal/notify" 34 "gitbay.org/gitbay/internal/notify"
35 "gitbay.org/gitbay/internal/packlimit"
35 "gitbay.org/gitbay/internal/push" 36 "gitbay.org/gitbay/internal/push"
36 "gitbay.org/gitbay/internal/seal" 37 "gitbay.org/gitbay/internal/seal"
37 "gitbay.org/gitbay/internal/sshd" 38 "gitbay.org/gitbay/internal/sshd"
@@ -228,11 +229,20 @@ func serveCmd() *cobra.Command {
228 return control.RepoDir(cfg.Server.Root, owner, name) 229 return control.RepoDir(cfg.Server.Root, owner, name)
229 }, buildinfo.String()).Run(whCtx) 230 }, buildinfo.String()).Run(whCtx)
230 231
232 // One pack-generation budget for SSH, smart HTTP and git://.
233 // Anonymous clients ("ip:" principals) share all but one
234 // slot, so an account can always get the last.
235 packMax, packPer, packQueue, packWait := cfg.Limits.PackLimits()
236 packs := packlimit.New(packMax, packPer, packQueue, packWait)
237 if packMax > 1 {
238 packs.CapClass("ip:", packMax-1)
239 }
240
231 errCh := make(chan error, 3) 241 errCh := make(chan error, 3)
232 var sshSrv *sshd.Server 242 var sshSrv *sshd.Server
233 var sshLn, gitLn net.Listener 243 var sshLn, gitLn net.Listener
234 if cfg.SSH.Mode == "embedded" { 244 if cfg.SSH.Mode == "embedded" {
235 srv, err := sshd.New(cfg, st) 245 srv, err := sshd.New(cfg, st, packs)
236 if err != nil { 246 if err != nil {
237 return err 247 return err
238 } 248 }
@@ -249,7 +259,7 @@ func serveCmd() *cobra.Command {
249 slog.Info("ssh handled by host sshd (ssh.mode = system)") 259 slog.Info("ssh handled by host sshd (ssh.mode = system)")
250 } 260 }
251 261
252 web := httpd.New(cfg, st) 262 web := httpd.New(cfg, st, packs)
253 // Header and idle timeouts bound what an idle or slow client can 263 // Header and idle timeouts bound what an idle or slow client can
254 // hold open. No write timeout: archives and upload-pack stream 264 // hold open. No write timeout: archives and upload-pack stream
255 // for as long as they take (#104). 265 // for as long as they take (#104).
@@ -338,7 +348,7 @@ func serveCmd() *cobra.Command {
338 } 348 }
339 slog.Info("git-daemon listening", "addr", gln.Addr()) 349 slog.Info("git-daemon listening", "addr", gln.Addr())
340 gitLn = gln 350 gitLn = gln
341 go func() { errCh <- gitd.New(cfg, st).Serve(gln) }() 351 go func() { errCh <- gitd.New(cfg, st, packs).Serve(gln) }()
342 } 352 }
343 353
344 select { 354 select {
cmd/gitbayd/system.go +3 −1
@@ -99,7 +99,9 @@ func shellCmd() *cobra.Command {
99 fmt.Fprintf(os.Stderr, "gitbay control plane: interactive shells are not available.\nTry: ssh <host> help\n") 99 fmt.Fprintf(os.Stderr, "gitbay control plane: interactive shells are not available.\nTry: ssh <host> help\n")
100 os.Exit(protocol.ExitUsage) 100 os.Exit(protocol.ExitUsage)
101 } 101 }
102 code := sshd.Exec(cfg, st, user, key, control.ParseTerm(os.Getenv("GITBAY_TERM")), cmdline, os.Stdin, os.Stdout, os.Stderr, nil, nil, nil) 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 st.Close() 105 st.Close()
104 os.Exit(code) 106 os.Exit(code)
105 return nil 107 return nil
deploy/clonebench.sh added +24
@@ -0,0 +1,24 @@
1#!/bin/sh
2# clonebench.sh <clone-url> <n>: start n full bare clones of <clone-url>
3# at once and print each one's wall time and outcome, then the total.
4# Run from a machine other than the server, against a public repository.
5set -eu
6url=$1
7n=$2
8dir=$(mktemp -d)
9trap 'rm -rf "$dir"' EXIT
10start=$(date +%s)
11i=1
12while [ "$i" -le "$n" ]; do
13 (
14 s=$(date +%s)
15 if git clone --quiet --bare "$url" "$dir/$i.git" 2>"$dir/$i.err"; then
16 echo "$i ok $(( $(date +%s) - s ))s"
17 else
18 echo "$i failed $(( $(date +%s) - s ))s: $(head -n 1 "$dir/$i.err")"
19 fi
20 ) &
21 i=$((i + 1))
22done
23wait
24echo "total $(( $(date +%s) - start ))s for $n clones"
internal/config/config.go +50
@@ -7,6 +7,7 @@ import (
7 "encoding/pem" 7 "encoding/pem"
8 "errors" 8 "errors"
9 "fmt" 9 "fmt"
10 "math"
10 "net" 11 "net"
11 "os" 12 "os"
12 "path/filepath" 13 "path/filepath"
@@ -24,6 +25,15 @@ import (
24// and webhook deliveries. 25// and webhook deliveries.
25const DefaultWriteRate = 60 26const DefaultWriteRate = 60
26 27
28// Pack generation defaults for a four-core host: a full clone of a large
29// repository runs git at about 1.5 cores (Performance wiki page).
30const (
31 DefaultPackConcurrency = 3
32 DefaultPackPerPrincipal = 2
33 DefaultPackQueue = 32
34 DefaultPackQueueWait = time.Minute
35)
36
27type Config struct { 37type Config struct {
28 Server Server `toml:"server"` 38 Server Server `toml:"server"`
29 SSH SSH `toml:"ssh"` 39 SSH SSH `toml:"ssh"`
@@ -214,6 +224,41 @@ type Limits struct {
214 // account. 224 // account.
215 MaxReposPerUser int `toml:"max_repos_per_user"` 225 MaxReposPerUser int `toml:"max_repos_per_user"`
216 MaxBytesPerUser int64 `toml:"max_bytes_per_user"` 226 MaxBytesPerUser int64 `toml:"max_bytes_per_user"`
227 // PackConcurrency caps git pack generation (upload-pack and
228 // upload-archive) running at once across SSH, smart HTTP and git://.
229 // PackPerPrincipal caps it per account, or per client address on the
230 // anonymous transports. PackQueue is how many may wait for a slot,
231 // for at most PackQueueWait ("60s"). For the three counts 0 takes the
232 // default and a negative value turns that bound off.
233 PackConcurrency int `toml:"pack_concurrency"`
234 PackPerPrincipal int `toml:"pack_per_principal"`
235 PackQueue int `toml:"pack_queue"`
236 PackQueueWait string `toml:"pack_queue_wait"`
237}
238
239// PackLimits resolves the pack_* settings for packlimit.New. A zero
240// max or per is no bound; an unbounded queue is math.MaxInt, since
241// packlimit reads a zero queue as no queue at all.
242func (l Limits) PackLimits() (max, per, queue int, wait time.Duration) {
243 pick := func(v, def int) int {
244 switch {
245 case v == 0:
246 return def
247 case v < 0:
248 return 0
249 }
250 return v
251 }
252 queue = pick(l.PackQueue, DefaultPackQueue)
253 if l.PackQueue < 0 {
254 queue = math.MaxInt
255 }
256 wait = DefaultPackQueueWait
257 if d, err := time.ParseDuration(l.PackQueueWait); err == nil && d > 0 {
258 wait = d
259 }
260 return pick(l.PackConcurrency, DefaultPackConcurrency),
261 pick(l.PackPerPrincipal, DefaultPackPerPrincipal), queue, wait
217} 262}
218 263
219type Mail struct { 264type Mail struct {
@@ -444,6 +489,11 @@ func (c Config) Validate() error {
444 if c.Limits.MaxReposPerUser < 0 || c.Limits.MaxBytesPerUser < 0 || c.Limits.MaxSnippetsPerUser < 0 { 489 if c.Limits.MaxReposPerUser < 0 || c.Limits.MaxBytesPerUser < 0 || c.Limits.MaxSnippetsPerUser < 0 {
445 errs = append(errs, errors.New("limits.max_repos_per_user, max_bytes_per_user and max_snippets_per_user must not be negative")) 490 errs = append(errs, errors.New("limits.max_repos_per_user, max_bytes_per_user and max_snippets_per_user must not be negative"))
446 } 491 }
492 if w := c.Limits.PackQueueWait; w != "" {
493 if d, err := time.ParseDuration(w); err != nil || d <= 0 {
494 errs = append(errs, fmt.Errorf("limits.pack_queue_wait %q must be a positive duration such as 60s", w))
495 }
496 }
447 if c.Push.Enabled { 497 if c.Push.Enabled {
448 for _, f := range []struct{ name, val string }{ 498 for _, f := range []struct{ name, val string }{
449 {"push.key_file", c.Push.KeyFile}, 499 {"push.key_file", c.Push.KeyFile},
internal/config/config_test.go +22
@@ -6,12 +6,14 @@ import (
6 "crypto/rand" 6 "crypto/rand"
7 "crypto/x509" 7 "crypto/x509"
8 "encoding/pem" 8 "encoding/pem"
9 "math"
9 "os" 10 "os"
10 "path/filepath" 11 "path/filepath"
11 "strings" 12 "strings"
12 "testing" 13 "testing"
13 14
14 "filippo.io/age" 15 "filippo.io/age"
16 "time"
15) 17)
16 18
17func writeConfig(t *testing.T, body string) string { 19func writeConfig(t *testing.T, body string) string {
@@ -46,12 +48,32 @@ func TestLoadMinimal(t *testing.T) {
46 } 48 }
47} 49}
48 50
51func TestPackLimits(t *testing.T) {
52 max, per, queue, wait := Limits{}.PackLimits()
53 if max != DefaultPackConcurrency || per != DefaultPackPerPrincipal || queue != DefaultPackQueue || wait != DefaultPackQueueWait {
54 t.Fatalf("defaults: %d %d %d %s", max, per, queue, wait)
55 }
56 max, per, queue, wait = Limits{PackConcurrency: -1, PackPerPrincipal: -1, PackQueue: -1, PackQueueWait: "5s"}.PackLimits()
57 if max != 0 || per != 0 || queue != math.MaxInt || wait != 5*time.Second {
58 t.Fatalf("off: %d %d %d %s", max, per, queue, wait)
59 }
60 max, per, queue, _ = Limits{PackConcurrency: 8, PackPerPrincipal: 3, PackQueue: 64}.PackLimits()
61 if max != 8 || per != 3 || queue != 64 {
62 t.Fatalf("set: %d %d %d", max, per, queue)
63 }
64}
65
49func TestContradictions(t *testing.T) { 66func TestContradictions(t *testing.T) {
50 cases := []struct { 67 cases := []struct {
51 name string 68 name string
52 body string 69 body string
53 wantErr string 70 wantErr string
54 }{ 71 }{
72 {
73 "bad pack_queue_wait",
74 minimal + "\n[limits]\npack_queue_wait = \"soon\"\n",
75 "limits.pack_queue_wait",
76 },
55 { 77 {
56 "registration open without smtp", 78 "registration open without smtp",
57 minimal + "\n[registration]\nmode = \"open\"\n", 79 minimal + "\n[registration]\nmode = \"open\"\n",
internal/gitd/gitd.go +43 −12
@@ -7,24 +7,26 @@ import (
7 "fmt" 7 "fmt"
8 "io" 8 "io"
9 "net" 9 "net"
10 "os"
11 "os/exec"
12 "strconv" 10 "strconv"
13 "strings" 11 "strings"
14 "time" 12 "time"
15 13
16 "gitbay.org/gitbay/internal/config" 14 "gitbay.org/gitbay/internal/config"
17 "gitbay.org/gitbay/internal/control" 15 "gitbay.org/gitbay/internal/control"
16 "gitbay.org/gitbay/internal/gitutil"
17 "gitbay.org/gitbay/internal/packlimit"
18 "gitbay.org/gitbay/internal/store" 18 "gitbay.org/gitbay/internal/store"
19 "gitbay.org/gitbay/internal/toolpath"
20) 19)
21 20
22type Server struct { 21type Server struct {
23 cfg config.Config 22 cfg config.Config
24 st *store.Store 23 st *store.Store
24 packs *packlimit.Limiter
25} 25}
26 26
27func New(cfg config.Config, st *store.Store) *Server { return &Server{cfg: cfg, st: st} } 27func New(cfg config.Config, st *store.Store, packs *packlimit.Limiter) *Server {
28 return &Server{cfg: cfg, st: st, packs: packs}
29}
28 30
29func (s *Server) Serve(ln net.Listener) error { 31func (s *Server) Serve(ln net.Listener) error {
30 for { 32 for {
@@ -68,13 +70,42 @@ func (s *Server) handle(conn net.Conn) {
68 return 70 return
69 } 71 }
70 72
73 // A nil done: a queued client that leaves, or a restart, does not
74 // end the wait; only the limiter's wait does.
75 p := principal(conn.RemoteAddr())
76 release, err := s.packs.Acquire(nil, p)
77 if err != nil {
78 s.packs.Refused("git", p, err)
79 writeErr(conn, err.Error())
80 return
81 }
82 // Deferred before git runs, so it fires after git has exited and
83 // been waited for.
84 defer release()
85 out, stalled, unwatch := s.packs.Watch(conn)
86 defer unwatch()
87 kill := make(chan struct{})
88 finished := make(chan struct{})
89 defer close(finished)
90 go func() {
91 select {
92 case <-finished:
93 case <-stalled:
94 close(kill)
95 // A write blocked on a client that stopped reading outlives
96 // git; closing the connection ends it and the stdin copy.
97 conn.Close()
98 }
99 }()
100
71 dir := control.RepoDir(s.cfg.Server.Root, repo.OwnerName, repo.Name) 101 dir := control.RepoDir(s.cfg.Server.Root, repo.OwnerName, repo.Name)
72 cmd := exec.Command(toolpath.Look("git"), "upload-pack", dir) 102 gitutil.Transport("git-upload-pack", dir, conn, out, io.Discard, protoEnv, 0, kill)
73 cmd.Env = append(os.Environ(), protoEnv...) 103}
74 cmd.Stdin = conn 104
75 cmd.Stdout = conn 105// principal is the pack-limit principal for a client at addr.
76 cmd.Stderr = io.Discard 106func principal(addr net.Addr) string {
77 cmd.Run() 107 host, _, _ := net.SplitHostPort(addr.String())
108 return packlimit.AddrPrincipal(host)
78} 109}
79 110
80func readPktLine(r io.Reader) (string, error) { 111func readPktLine(r io.Reader) (string, error) {
internal/gitd/gitd_test.go added +62
@@ -0,0 +1,62 @@
1package gitd
2
3import (
4 "fmt"
5 "net"
6 "path/filepath"
7 "strings"
8 "testing"
9 "time"
10
11 "gitbay.org/gitbay/internal/config"
12 "gitbay.org/gitbay/internal/packlimit"
13 "gitbay.org/gitbay/internal/store"
14)
15
16func TestBusyAnswersERR(t *testing.T) {
17 st, err := store.Open(filepath.Join(t.TempDir(), "gitbay.db"))
18 if err != nil {
19 t.Fatal(err)
20 }
21 defer st.Close()
22 if err := st.MigrateUp(); err != nil {
23 t.Fatal(err)
24 }
25 uid, err := st.CreateUser("alice", false)
26 if err != nil {
27 t.Fatal(err)
28 }
29 repoID, err := st.CreateRepo("user", uid, "app", "public")
30 if err != nil {
31 t.Fatal(err)
32 }
33 if _, err := st.UpdateRepoSettings(repoID, func(rs *store.RepoSettings) { rs.GitDaemon = true }); err != nil {
34 t.Fatal(err)
35 }
36 packs := packlimit.New(1, 0, 0, time.Second)
37 hold, _ := packs.Acquire(nil, "ip:elsewhere")
38 defer hold()
39
40 s := New(config.Config{Server: config.Server{Root: t.TempDir()}}, st, packs)
41 client, server := net.Pipe()
42 defer client.Close()
43 go s.handle(server)
44 req := "git-upload-pack /alice/app.git\x00host=x\x00"
45 fmt.Fprintf(client, "%04x%s", len(req)+4, req)
46 client.SetReadDeadline(time.Now().Add(5 * time.Second))
47 line, err := readPktLine(client)
48 if err != nil || !strings.HasPrefix(line, "ERR ") || !strings.Contains(line, "busy") {
49 t.Fatalf("got %q, %v", line, err)
50 }
51}
52
53func TestPrincipal(t *testing.T) {
54 for addr, want := range map[net.Addr]string{
55 &net.TCPAddr{IP: net.ParseIP("192.0.2.7"), Port: 9418}: "ip:192.0.2.7",
56 &net.TCPAddr{IP: net.ParseIP("2001:db8:1:2:3:4:5:6"), Port: 9418}: "ip:2001:db8:1:2::/64",
57 } {
58 if got := principal(addr); got != want {
59 t.Errorf("principal(%v) = %q, want %q", addr, got, want)
60 }
61 }
62}
internal/gitutil/gitutil.go +12
@@ -45,6 +45,11 @@ func Transport(service, repoPath string, stdin io.Reader, stdout, errW io.Writer
45 if service == "git-receive-pack" && maxPack > 0 { 45 if service == "git-receive-pack" && maxPack > 0 {
46 args = []string{"-c", fmt.Sprintf("receive.maxInputSize=%d", maxPack)} 46 args = []string{"-c", fmt.Sprintf("receive.maxInputSize=%d", maxPack)}
47 } 47 }
48 if service == "git-upload-pack" {
49 // Keepalives while pack-objects is still counting keep a
50 // healthy clone writing; a limited transport kills one that goes quiet.
51 args = []string{"-c", "uploadpack.keepAlive=5"}
52 }
48 args = append(args, strings.TrimPrefix(service, "git-"), repoPath) 53 args = append(args, strings.TrimPrefix(service, "git-"), repoPath)
49 default: 54 default:
50 return fmt.Errorf("unknown service %q", service) 55 return fmt.Errorf("unknown service %q", service)
@@ -54,6 +59,13 @@ func Transport(service, repoPath string, stdin io.Reader, stdout, errW io.Writer
54 cmd.Stdin = stdin 59 cmd.Stdin = stdin
55 cmd.Stdout = stdout 60 cmd.Stdout = stdout
56 cmd.Stderr = errW 61 cmd.Stderr = errW
62 return RunUntil(cmd, cancel)
63}
64
65// RunUntil runs cmd in its own process group. Closing cancel kills the
66// group; RunUntil returns only once cmd has been waited for. A nil
67// cancel never fires.
68func RunUntil(cmd *exec.Cmd, cancel <-chan struct{}) error {
57 ownProcessGroup(cmd) 69 ownProcessGroup(cmd)
58 if err := cmd.Start(); err != nil { 70 if err := cmd.Start(); err != nil {
59 return err 71 return err
internal/gitutil/read.go +6 −1
@@ -132,12 +132,17 @@ var ErrArchiveTooLarge = errors.New("archive exceeds the size limit")
132// Archive streams a tar.gz of ref to w, within archiveTimeout and 132// Archive streams a tar.gz of ref to w, within archiveTimeout and
133// MaxArchiveBytes. Past either, git is killed and the error says which. 133// MaxArchiveBytes. Past either, git is killed and the error says which.
134func Archive(dir, ref, prefix string, w io.Writer) error { 134func Archive(dir, ref, prefix string, w io.Writer) error {
135 return ArchiveUntil(dir, ref, prefix, w, nil)
136}
137
138// ArchiveUntil is Archive, killing git when stop closes.
139func ArchiveUntil(dir, ref, prefix string, w io.Writer, stop <-chan struct{}) error {
135 ctx, cancel := context.WithTimeout(context.Background(), archiveTimeout) 140 ctx, cancel := context.WithTimeout(context.Background(), archiveTimeout)
136 defer cancel() 141 defer cancel()
137 cmd := exec.CommandContext(ctx, toolpath.Look("git"), "-C", dir, "archive", "--format=tar.gz", "--prefix="+prefix+"/", "--end-of-options", ref) 142 cmd := exec.CommandContext(ctx, toolpath.Look("git"), "-C", dir, "archive", "--format=tar.gz", "--prefix="+prefix+"/", "--end-of-options", ref)
138 lw := &cappedWriter{w: w, left: MaxArchiveBytes, stop: cancel} 143 lw := &cappedWriter{w: w, left: MaxArchiveBytes, stop: cancel}
139 cmd.Stdout = lw 144 cmd.Stdout = lw
140 err := cmd.Run() 145 err := RunUntil(cmd, stop)
141 switch { 146 switch {
142 case lw.exceeded: 147 case lw.exceeded:
143 return ErrArchiveTooLarge 148 return ErrArchiveTooLarge
internal/hookd/hookd.go +38 −10
@@ -122,8 +122,10 @@ func (s *Server) handle(conn net.Conn) {
122 defer conn.Close() 122 defer conn.Close()
123 dec := json.NewDecoder(conn) 123 dec := json.NewDecoder(conn)
124 enc := json.NewEncoder(conn) 124 enc := json.NewEncoder(conn)
125 if err := checkPeer(conn); err != nil { 125 if err := peerCheck(conn); err != nil {
126 slog.Warn("hook socket: refused connection", "err", err) 126 slog.Warn("hook socket: refused connection", "err", err)
127 // Nothing about the request is known yet, and no account.
128 control.AuditRefused(s.st, 0, "refused hook", map[string]any{"reason": err.Error()})
127 enc.Encode(Response{Allow: false, Message: "hook socket: " + err.Error()}) 129 enc.Encode(Response{Allow: false, Message: "hook socket: " + err.Error()})
128 return 130 return
129 } 131 }
@@ -132,7 +134,9 @@ func (s *Server) handle(conn net.Conn) {
132 enc.Encode(Response{Allow: false, Message: "bad hook request"}) 134 enc.Encode(Response{Allow: false, Message: "bad hook request"})
133 return 135 return
134 } 136 }
135 if msg := s.authorize(req); msg != "" { 137 if actor, msg := s.authorize(req); msg != "" {
138 control.AuditRefused(s.st, actor, "refused hook",
139 map[string]any{"repo_id": req.RepoID, "hook": req.Hook, "reason": msg})
136 enc.Encode(Response{Allow: false, Message: msg}) 140 enc.Encode(Response{Allow: false, Message: msg})
137 return 141 return
138 } 142 }
@@ -149,18 +153,42 @@ func (s *Server) handle(conn net.Conn) {
149 153
150// authorize ties a request to a receive-pack sshd started: its token 154// authorize ties a request to a receive-pack sshd started: its token
151// must be live and name the same repository, account and key scope. 155// must be live and name the same repository, account and key scope.
152func (s *Server) authorize(req Request) string { 156// On a refusal actor is the token's account when the token is live,
157// and 0 otherwise: the request's own user id is only a claim.
158func (s *Server) authorize(req Request) (actor int64, msg string) {
153 if req.Token == "" { 159 if req.Token == "" {
154 return "push not started by this server" 160 return 0, "push not started by this server"
155 } 161 }
156 tok, err := s.st.PushTokenByHash(store.HashToken(req.Token)) 162 tok, err := s.st.PushTokenByHash(store.HashToken(req.Token))
157 if err != nil { 163 if err != nil {
158 return "push not started by this server" 164 return 0, "push not started by this server"
159 } 165 }
160 if tok.RepoID != req.RepoID || tok.UserID != req.UserID || tok.Scope != req.Scope { 166 if tok.RepoID != req.RepoID || tok.UserID != req.UserID || tok.Scope != req.Scope {
161 return "push token does not match this request" 167 return tok.UserID, "push token does not match this request"
162 } 168 }
163 return "" 169 return 0, ""
170}
171
172// peerCheck is checkPeer; tests replace it.
173var peerCheck = checkPeer
174
175// auditedRefs is how many ref names a refused-push row keeps; the rest
176// are counted, so one push of many refs cannot write an unbounded row.
177const auditedRefs = 20
178
179// refusePush answers a pre-receive refusal and audits it.
180func (s *Server) refusePush(enc *json.Encoder, req Request, repo store.Repo, msg string) {
181 n := min(len(req.Updates), auditedRefs)
182 refs := make([]string, n)
183 for i, u := range req.Updates[:n] {
184 refs[i] = u.Ref
185 }
186 data := map[string]any{"repo": repo.Path(), "refs": refs, "reason": msg}
187 if more := len(req.Updates) - n; more > 0 {
188 data["more_refs"] = more
189 }
190 control.AuditRefused(s.st, req.UserID, "refused push", data)
191 enc.Encode(Response{Allow: false, Message: msg})
164} 192}
165 193
166func (s *Server) preReceive(req Request, dec *json.Decoder, enc *json.Encoder) { 194func (s *Server) preReceive(req Request, dec *json.Decoder, enc *json.Encoder) {
@@ -170,11 +198,11 @@ func (s *Server) preReceive(req Request, dec *json.Decoder, enc *json.Encoder) {
170 return 198 return
171 } 199 }
172 if msg := policy.CheckPush(repo, req.Updates); msg != "" { 200 if msg := policy.CheckPush(repo, req.Updates); msg != "" {
173 enc.Encode(Response{Allow: false, Message: msg}) 201 s.refusePush(enc, req, repo, msg)
174 return 202 return
175 } 203 }
176 if msg := s.releaseAnchors(repo, req.Updates); msg != "" { 204 if msg := s.releaseAnchors(repo, req.Updates); msg != "" {
177 enc.Encode(Response{Allow: false, Message: msg}) 205 s.refusePush(enc, req, repo, msg)
178 return 206 return
179 } 207 }
180 if !repo.Settings.RequireSignedCommits { 208 if !repo.Settings.RequireSignedCommits {
@@ -220,7 +248,7 @@ func (s *Server) preReceive(req Request, dec *json.Decoder, enc *json.Encoder) {
220 } 248 }
221 } 249 }
222 if refusal != "" { 250 if refusal != "" {
223 enc.Encode(Response{Allow: false, Message: refusal}) 251 s.refusePush(enc, req, repo, refusal)
224 return 252 return
225 } 253 }
226 enc.Encode(Response{Allow: true}) 254 enc.Encode(Response{Allow: true})
internal/hookd/socket_test.go +92
@@ -1,12 +1,17 @@
1package hookd 1package hookd
2 2
3import ( 3import (
4 "encoding/json"
5 "errors"
6 "fmt"
7 "net"
4 "os" 8 "os"
5 "path/filepath" 9 "path/filepath"
6 "strings" 10 "strings"
7 "testing" 11 "testing"
8 12
9 "gitbay.org/gitbay/internal/config" 13 "gitbay.org/gitbay/internal/config"
14 "gitbay.org/gitbay/internal/policy"
10 "gitbay.org/gitbay/internal/store" 15 "gitbay.org/gitbay/internal/store"
11) 16)
12 17
@@ -103,3 +108,90 @@ func TestHookRequestNeedsItsPushToken(t *testing.T) {
103 t.Fatalf("finished push: %+v, %v", resp, err) 108 t.Fatalf("finished push: %+v, %v", resp, err)
104 } 109 }
105} 110}
111
112func refusedRows(t *testing.T, st *store.Store, action string) []store.AuditEntry {
113 t.Helper()
114 rows, err := st.AuditEntries(store.AuditFilter{ActionPrefix: action, Limit: 10})
115 if err != nil {
116 t.Fatal(err)
117 }
118 return rows
119}
120
121// Refused hook requests and refused pushes are audited; the token never
122// lands in a row (#275).
123func TestHookRefusalsAreAudited(t *testing.T) {
124 sock, st, repoID, uid := serveSocket(t)
125
126 forged := Request{Hook: "pre-receive", RepoID: repoID, UserID: uid, Scope: "full", Token: "not-a-live-token"}
127 if resp, err := Ask(sock, forged, nil); err != nil || resp.Allow {
128 t.Fatalf("forged: %+v, %v", resp, err)
129 }
130 rows := refusedRows(t, st, "refused hook")
131 if len(rows) != 1 || rows[0].Actor != "" || strings.Contains(rows[0].Data, forged.Token) ||
132 !strings.Contains(rows[0].Data, "not started by this server") || !strings.Contains(rows[0].Data, `"hook":"pre-receive"`) {
133 t.Fatalf("refused hook rows: %+v", rows)
134 }
135
136 token, err := st.CreatePushToken(repoID, uid, "full")
137 if err != nil {
138 t.Fatal(err)
139 }
140 req := Request{Hook: "pre-receive", RepoID: repoID, UserID: uid, Scope: "full", Token: token,
141 Updates: []policy.RefUpdate{{Ref: "refs/merge-requests/1/head", Old: zeroSHA40, New: strings.Repeat("a", 40)}}}
142 if resp, err := Ask(sock, req, nil); err != nil || resp.Allow {
143 t.Fatalf("push to a server-owned ref: %+v, %v", resp, err)
144 }
145 rows = refusedRows(t, st, "refused push")
146 if len(rows) != 1 || rows[0].Actor != "alice" || strings.Contains(rows[0].Data, token) ||
147 !strings.Contains(rows[0].Data, "alice/app") || !strings.Contains(rows[0].Data, "refs/merge-requests/1/head") {
148 t.Fatalf("refused push rows: %+v", rows)
149 }
150}
151
152// A connection from another uid is audited with no actor.
153func TestPeerRefusalIsAudited(t *testing.T) {
154 old := peerCheck
155 peerCheck = func(net.Conn) error { return errors.New("peer uid not permitted") }
156 t.Cleanup(func() { peerCheck = old })
157 sock, st, repoID, uid := serveSocket(t)
158 if resp, err := Ask(sock, Request{Hook: "pre-receive", RepoID: repoID, UserID: uid}, nil); err != nil || resp.Allow {
159 t.Fatalf("refused peer: %+v, %v", resp, err)
160 }
161 rows := refusedRows(t, st, "refused hook")
162 if len(rows) != 1 || rows[0].Actor != "" || !strings.Contains(rows[0].Data, "peer uid not permitted") {
163 t.Fatalf("refused hook rows: %+v", rows)
164 }
165}
166
167// A refused push of many refs records the first auditedRefs names and
168// a count of the rest.
169func TestRefusedPushCapsRefs(t *testing.T) {
170 sock, st, repoID, uid := serveSocket(t)
171 token, err := st.CreatePushToken(repoID, uid, "full")
172 if err != nil {
173 t.Fatal(err)
174 }
175 req := Request{Hook: "pre-receive", RepoID: repoID, UserID: uid, Scope: "full", Token: token}
176 for i := range 500 {
177 req.Updates = append(req.Updates, policy.RefUpdate{Ref: fmt.Sprintf("refs/merge-requests/%d/head", i),
178 Old: zeroSHA40, New: strings.Repeat("a", 40)})
179 }
180 if resp, err := Ask(sock, req, nil); err != nil || resp.Allow {
181 t.Fatalf("push: %+v, %v", resp, err)
182 }
183 rows := refusedRows(t, st, "refused push")
184 if len(rows) != 1 {
185 t.Fatalf("rows: %+v", rows)
186 }
187 var data struct {
188 Refs []string `json:"refs"`
189 MoreRefs int `json:"more_refs"`
190 }
191 if err := json.Unmarshal([]byte(rows[0].Data), &data); err != nil {
192 t.Fatal(err)
193 }
194 if len(data.Refs) != auditedRefs || data.MoreRefs != 500-auditedRefs || data.Refs[0] != "refs/merge-requests/0/head" {
195 t.Fatalf("refs %d, more %d", len(data.Refs), data.MoreRefs)
196 }
197}
internal/httpd/account_test.go +6 −6
@@ -34,7 +34,7 @@ func TestAccountPagePushToggleAndDevices(t *testing.T) {
34 t.Fatal(err) 34 t.Fatal(err)
35 } 35 }
36 36
37 s := New(config.Default(), st) 37 s := New(config.Default(), st, nil)
38 rr := httptest.NewRecorder() 38 rr := httptest.NewRecorder()
39 req := httptest.NewRequest("GET", "/settings", nil) 39 req := httptest.NewRequest("GET", "/settings", nil)
40 s.accountPage(rr, req, store.User{ID: uid, Username: "alice"}) 40 s.accountPage(rr, req, store.User{ID: uid, Username: "alice"})
@@ -78,7 +78,7 @@ func TestAccountSubmitNotifyPush(t *testing.T) {
78 t.Fatal(err) 78 t.Fatal(err)
79 } 79 }
80 u := store.User{ID: uid, Username: "alice"} 80 u := store.User{ID: uid, Username: "alice"}
81 s := New(config.Default(), st) 81 s := New(config.Default(), st, nil)
82 82
83 rr := submitAccountForm(t, s, u, url.Values{"field": {"notify-push"}, "push": {"on"}}) 83 rr := submitAccountForm(t, s, u, url.Values{"field": {"notify-push"}, "push": {"on"}})
84 if rr.Code != http.StatusSeeOther { 84 if rr.Code != http.StatusSeeOther {
@@ -117,7 +117,7 @@ func TestAccountSubmitDeviceRemove(t *testing.T) {
117 if err != nil { 117 if err != nil {
118 t.Fatal(err) 118 t.Fatal(err)
119 } 119 }
120 s := New(config.Default(), st) 120 s := New(config.Default(), st, nil)
121 121
122 idStr := strconv.FormatInt(id, 10) 122 idStr := strconv.FormatInt(id, 10)
123 123
@@ -158,7 +158,7 @@ func TestAccountPageMasksAShortDeviceToken(t *testing.T) {
158 t.Fatal(err) 158 t.Fatal(err)
159 } 159 }
160 160
161 s := New(config.Default(), st) 161 s := New(config.Default(), st, nil)
162 rr := httptest.NewRecorder() 162 rr := httptest.NewRecorder()
163 s.accountPage(rr, httptest.NewRequest("GET", "/settings", nil), store.User{ID: uid, Username: "alice"}) 163 s.accountPage(rr, httptest.NewRequest("GET", "/settings", nil), store.User{ID: uid, Username: "alice"})
164 164
@@ -212,7 +212,7 @@ func TestPinToggleDispatchesRepoPin(t *testing.T) {
212 t.Fatal(err) 212 t.Fatal(err)
213 } 213 }
214 214
215 s := New(config.Default(), st) 215 s := New(config.Default(), st, nil)
216 req := httptest.NewRequest("POST", "/alice/app/pin", nil) 216 req := httptest.NewRequest("POST", "/alice/app/pin", nil)
217 req.SetPathValue("owner", "alice") 217 req.SetPathValue("owner", "alice")
218 req.SetPathValue("repo", "app") 218 req.SetPathValue("repo", "app")
@@ -261,7 +261,7 @@ func TestWatchToggleCyclesThroughMuted(t *testing.T) {
261 t.Fatal(err) 261 t.Fatal(err)
262 } 262 }
263 263
264 s := New(config.Default(), st) 264 s := New(config.Default(), st, nil)
265 req := httptest.NewRequest("POST", "/alice/app/watch", nil) 265 req := httptest.NewRequest("POST", "/alice/app/watch", nil)
266 req.SetPathValue("owner", "alice") 266 req.SetPathValue("owner", "alice")
267 req.SetPathValue("repo", "app") 267 req.SetPathValue("repo", "app")
internal/httpd/admin_test.go +1 −1
@@ -48,7 +48,7 @@ func TestAdminPageShowsThePushQueue(t *testing.T) {
48 t.Fatal(err) 48 t.Fatal(err)
49 } 49 }
50 50
51 s := New(config.Default(), st) 51 s := New(config.Default(), st, nil)
52 rr := httptest.NewRecorder() 52 rr := httptest.NewRecorder()
53 s.adminPage(rr, httptest.NewRequest("GET", "/admin", nil), store.User{ID: uid, Username: "root", IsAdmin: true}) 53 s.adminPage(rr, httptest.NewRequest("GET", "/admin", nil), store.User{ID: uid, Username: "root", IsAdmin: true})
54 if rr.Code != http.StatusOK { 54 if rr.Code != http.StatusOK {
internal/httpd/anchors_test.go +2 −2
@@ -26,7 +26,7 @@ func TestMarkdownHeadingAnchors(t *testing.T) {
26// The stylesheet carries an ETag and a cache lifetime; a revalidation 26// The stylesheet carries an ETag and a cache lifetime; a revalidation
27// with the same tag is a 304 with no body (#132). 27// with the same tag is a 304 with no body (#132).
28func TestStylesheetRevalidates(t *testing.T) { 28func TestStylesheetRevalidates(t *testing.T) {
29 s := New(config.Default(), nil) 29 s := New(config.Default(), nil, nil)
30 first := httptest.NewRecorder() 30 first := httptest.NewRecorder()
31 s.stylesheet(first, httptest.NewRequest("GET", "/static/style.css", nil)) 31 s.stylesheet(first, httptest.NewRequest("GET", "/static/style.css", nil))
32 tag := first.Header().Get("ETag") 32 tag := first.Header().Get("ETag")
@@ -69,7 +69,7 @@ func TestStylesheetURLCarriesTheBuildHash(t *testing.T) {
69 t.Errorf("the page does not link %s", want) 69 t.Errorf("the page does not link %s", want)
70 } 70 }
71 71
72 s := New(config.Default(), nil) 72 s := New(config.Default(), nil, nil)
73 versioned := httptest.NewRecorder() 73 versioned := httptest.NewRecorder()
74 s.stylesheet(versioned, httptest.NewRequest("GET", "/static/style.css?v="+stylesheetHash, nil)) 74 s.stylesheet(versioned, httptest.NewRequest("GET", "/static/style.css?v="+stylesheetHash, nil))
75 if cc := versioned.Header().Get("Cache-Control"); !strings.Contains(cc, "immutable") { 75 if cc := versioned.Header().Get("Cache-Control"); !strings.Contains(cc, "immutable") {
internal/httpd/checkorigin_test.go +1 −1
@@ -40,7 +40,7 @@ func TestMutatingRoutesRequireCheckOrigin(t *testing.T) {
40 // there is no second code path in routes.go for looping over both to 40 // there is no second code path in routes.go for looping over both to
41 // reach; "open" alone matches production and is enough. 41 // reach; "open" alone matches production and is enough.
42 cfg.Registration.Mode = "open" 42 cfg.Registration.Mode = "open"
43 s := New(cfg, nil) 43 s := New(cfg, nil, nil)
44 44
45 for _, r := range s.Routes() { 45 for _, r := range s.Routes() {
46 if !r.Mutating { 46 if !r.Mutating {
internal/httpd/clientip_test.go +1 −1
@@ -28,7 +28,7 @@ func TestClientIPBehindProxy(t *testing.T) {
28 for _, tc := range cases { 28 for _, tc := range cases {
29 cfg := config.Default() 29 cfg := config.Default()
30 cfg.HTTP.TrustedProxies = tc.proxies 30 cfg.HTTP.TrustedProxies = tc.proxies
31 s := New(cfg, nil) 31 s := New(cfg, nil, nil)
32 r := httptest.NewRequest("GET", "/api/v1/read", nil) 32 r := httptest.NewRequest("GET", "/api/v1/read", nil)
33 r.RemoteAddr = tc.remote 33 r.RemoteAddr = tc.remote
34 if tc.xff != "" { 34 if tc.xff != "" {
internal/httpd/fonts_test.go +2 −2
@@ -17,7 +17,7 @@ import (
17// and the @font-face URLs were once maintained by hand and drifted, so 17// and the @font-face URLs were once maintained by hand and drifted, so
18// gitbay.org served no web font at all (#102). 18// gitbay.org served no web font at all (#102).
19func TestStylesheetFontsAreServed(t *testing.T) { 19func TestStylesheetFontsAreServed(t *testing.T) {
20 s := New(config.Default(), nil) 20 s := New(config.Default(), nil, nil)
21 byPattern := map[string]http.HandlerFunc{} 21 byPattern := map[string]http.HandlerFunc{}
22 for _, r := range s.Routes() { 22 for _, r := range s.Routes() {
23 if r.Method == "GET" { 23 if r.Method == "GET" {
@@ -49,7 +49,7 @@ func TestStylesheetFontsAreServed(t *testing.T) {
49// TestLandingImagesAreServed: every file under static/img has a route 49// TestLandingImagesAreServed: every file under static/img has a route
50// that answers 200 with an image or video type, and a video answers Range. 50// that answers 200 with an image or video type, and a video answers Range.
51func TestLandingImagesAreServed(t *testing.T) { 51func TestLandingImagesAreServed(t *testing.T) {
52 s := New(config.Default(), nil) 52 s := New(config.Default(), nil, nil)
53 byPattern := map[string]http.HandlerFunc{} 53 byPattern := map[string]http.HandlerFunc{}
54 for _, r := range s.Routes() { 54 for _, r := range s.Routes() {
55 if r.Method == "GET" { 55 if r.Method == "GET" {
internal/httpd/issuecreate_test.go +5 −5
@@ -79,7 +79,7 @@ func TestIssueCreateFormHasMilestoneAndAssigneeForWriter(t *testing.T) {
79 79
80 cfg := config.Default() 80 cfg := config.Default()
81 cfg.Web.Mode = "accounts" 81 cfg.Web.Mode = "accounts"
82 s := New(cfg, st) 82 s := New(cfg, st, nil)
83 req := httptest.NewRequest("GET", "/alice/app/issues/new", nil) 83 req := httptest.NewRequest("GET", "/alice/app/issues/new", nil)
84 req.SetPathValue("owner", "alice") 84 req.SetPathValue("owner", "alice")
85 req.SetPathValue("repo", "app") 85 req.SetPathValue("repo", "app")
@@ -125,7 +125,7 @@ func TestIssueCreateFormHidesMilestoneAndAssigneeForReader(t *testing.T) {
125 125
126 cfg := config.Default() 126 cfg := config.Default()
127 cfg.Web.Mode = "accounts" 127 cfg.Web.Mode = "accounts"
128 s := New(cfg, st) 128 s := New(cfg, st, nil)
129 req := httptest.NewRequest("GET", "/alice/app/issues/new", nil) 129 req := httptest.NewRequest("GET", "/alice/app/issues/new", nil)
130 req.SetPathValue("owner", "alice") 130 req.SetPathValue("owner", "alice")
131 req.SetPathValue("repo", "app") 131 req.SetPathValue("repo", "app")
@@ -198,7 +198,7 @@ func TestIssueCreateSubmitReaderLabelIsDropped(t *testing.T) {
198 t.Fatal(err) 198 t.Fatal(err)
199 } 199 }
200 200
201 s := New(config.Default(), st) 201 s := New(config.Default(), st, nil)
202 form := url.Values{ 202 form := url.Values{
203 "title": {"a bug"}, 203 "title": {"a bug"},
204 "body": {"steps"}, 204 "body": {"steps"},
@@ -254,7 +254,7 @@ func TestIssueCreateSubmitSetsMilestoneAndAssignee(t *testing.T) {
254 t.Fatal(err) 254 t.Fatal(err)
255 } 255 }
256 256
257 s := New(config.Default(), st) 257 s := New(config.Default(), st, nil)
258 form := url.Values{ 258 form := url.Values{
259 "title": {"needs a fix"}, 259 "title": {"needs a fix"},
260 "body": {"details"}, 260 "body": {"details"},
@@ -308,7 +308,7 @@ func TestIssueCreateSubmitBadAssigneeCreatesNothing(t *testing.T) {
308 t.Fatal(err) 308 t.Fatal(err)
309 } 309 }
310 310
311 s := New(config.Default(), st) 311 s := New(config.Default(), st, nil)
312 form := url.Values{ 312 form := url.Values{
313 "title": {"needs a fix"}, 313 "title": {"needs a fix"},
314 "body": {"details"}, 314 "body": {"details"},
internal/httpd/logincookie_test.go +1 −1
@@ -52,7 +52,7 @@ func TestLoginNoStoreHeader(t *testing.T) {
52 if err := st.MigrateUp(); err != nil { 52 if err := st.MigrateUp(); err != nil {
53 t.Fatal(err) 53 t.Fatal(err)
54 } 54 }
55 s := New(config.Default(), st) 55 s := New(config.Default(), st, nil)
56 rr := httptest.NewRecorder() 56 rr := httptest.NewRecorder()
57 req := httptest.NewRequest("GET", "/login?token=bogus", nil) 57 req := httptest.NewRequest("GET", "/login?token=bogus", nil)
58 s.login(rr, req) 58 s.login(rr, req)
internal/httpd/logindisabled_test.go +1 −1
@@ -40,7 +40,7 @@ func TestLoginRefusesTokenForDisabledAccount(t *testing.T) {
40 t.Fatal(err) 40 t.Fatal(err)
41 } 41 }
42 42
43 s := New(config.Default(), st) 43 s := New(config.Default(), st, nil)
44 rr := httptest.NewRecorder() 44 rr := httptest.NewRecorder()
45 req := httptest.NewRequest("GET", "/login?token="+tok, nil) 45 req := httptest.NewRequest("GET", "/login?token="+tok, nil)
46 s.login(rr, req) 46 s.login(rr, req)
internal/httpd/mrrangediff_test.go +3 −3
@@ -36,7 +36,7 @@ func TestMRRangeDiffPageRendersCommandOutput(t *testing.T) {
36 t.Fatal(err) 36 t.Fatal(err)
37 } 37 }
38 38
39 s := New(config.Default(), st) 39 s := New(config.Default(), st, nil)
40 req := httptest.NewRequest("GET", "/alice/app/mrs/1/range-diff", nil) 40 req := httptest.NewRequest("GET", "/alice/app/mrs/1/range-diff", nil)
41 req.SetPathValue("owner", "alice") 41 req.SetPathValue("owner", "alice")
42 req.SetPathValue("repo", "app") 42 req.SetPathValue("repo", "app")
@@ -180,7 +180,7 @@ func loginCookie(t *testing.T, st *store.Store, userID int64) *http.Cookie {
180// user with no access (#269). 180// user with no access (#269).
181func TestMRRangeDiffPagePrivateRepo(t *testing.T) { 181func TestMRRangeDiffPagePrivateRepo(t *testing.T) {
182 st, cfg, alice, bob, repo, n, _, _, title := rangeDiffFixture(t) 182 st, cfg, alice, bob, repo, n, _, _, title := rangeDiffFixture(t)
183 s := New(cfg, st) 183 s := New(cfg, st, nil)
184 184
185 newReq := func(cookie *http.Cookie) (*httptest.ResponseRecorder, *http.Request) { 185 newReq := func(cookie *http.Cookie) (*httptest.ResponseRecorder, *http.Request) {
186 req := httptest.NewRequest("GET", "/alice/secret/mrs/"+strconv.FormatInt(n, 10)+"/range-diff", nil) 186 req := httptest.NewRequest("GET", "/alice/secret/mrs/"+strconv.FormatInt(n, 10)+"/range-diff", nil)
@@ -232,7 +232,7 @@ func TestMRRangeDiffPagePrivateRepo(t *testing.T) {
232// range-diff between exactly those two. 232// range-diff between exactly those two.
233func TestMRRangeDiffPageFromToQuery(t *testing.T) { 233func TestMRRangeDiffPageFromToQuery(t *testing.T) {
234 st, cfg, alice, _, repo, n, v1, v2, _ := rangeDiffFixture(t) 234 st, cfg, alice, _, repo, n, v1, v2, _ := rangeDiffFixture(t)
235 s := New(cfg, st) 235 s := New(cfg, st, nil)
236 cookie := loginCookie(t, st, alice.ID) 236 cookie := loginCookie(t, st, alice.ID)
237 237
238 newReq := func(query string) (*httptest.ResponseRecorder, *http.Request) { 238 newReq := func(query string) (*httptest.ResponseRecorder, *http.Request) {
internal/httpd/mrslist_test.go +1 −1
@@ -37,7 +37,7 @@ func TestMRsListContributionHintByAccess(t *testing.T) {
37 37
38 cfg := config.Default() 38 cfg := config.Default()
39 cfg.Web.Mode = "accounts" 39 cfg.Web.Mode = "accounts"
40 s := New(cfg, st) 40 s := New(cfg, st, nil)
41 41
42 // mrs reads the viewer through s.viewer(r), which resolves a 42 // mrs reads the viewer through s.viewer(r), which resolves a
43 // session cookie (internal/httpd/accounts.go:37-47) rather than 43 // session cookie (internal/httpd/accounts.go:37-47) rather than
internal/httpd/packlimit_test.go added +348
@@ -0,0 +1,348 @@
1package httpd
2
3import (
4 "context"
5 "crypto/rand"
6 "fmt"
7 "net"
8 "net/http"
9 "net/http/httptest"
10 "os"
11 "os/exec"
12 "path/filepath"
13 "strconv"
14 "strings"
15 "sync"
16 "testing"
17 "time"
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
26func busyServer(t *testing.T) *Server {
27 t.Helper()
28 st, err := store.Open(filepath.Join(t.TempDir(), "gitbay.db"))
29 if err != nil {
30 t.Fatal(err)
31 }
32 t.Cleanup(func() { st.Close() })
33 if err := st.MigrateUp(); err != nil {
34 t.Fatal(err)
35 }
36 uid, err := st.CreateUser("alice", false)
37 if err != nil {
38 t.Fatal(err)
39 }
40 if _, err := st.CreateRepo("user", uid, "app", "public"); err != nil {
41 t.Fatal(err)
42 }
43 packs := packlimit.New(1, 0, 0, time.Second)
44 hold, err := packs.Acquire(nil, "ip:elsewhere")
45 if err != nil {
46 t.Fatal(err)
47 }
48 t.Cleanup(hold)
49 var cfg config.Config
50 cfg.Server.Root = t.TempDir()
51 return &Server{cfg: cfg, st: st, packs: packs, stopping: make(chan struct{})}
52}
53
54func post(s *Server, body string) *httptest.ResponseRecorder {
55 r := httptest.NewRequest("POST", "/alice/app/git-upload-pack", strings.NewReader(body))
56 r.SetPathValue("owner", "alice")
57 r.SetPathValue("repo", "app")
58 w := httptest.NewRecorder()
59 s.uploadPack(w, r)
60 return w
61}
62
63func TestUploadPackBusyIs503(t *testing.T) {
64 w := post(busyServer(t), "0000")
65 if w.Code != http.StatusServiceUnavailable || w.Header().Get("Retry-After") == "" {
66 t.Fatalf("status %d, Retry-After %q", w.Code, w.Header().Get("Retry-After"))
67 }
68}
69
70// A protocol v2 ref listing generates no pack and is never queued.
71func TestLsRefsBypassesTheLimit(t *testing.T) {
72 w := post(busyServer(t), "0014command=ls-refs\n0000")
73 if w.Code == http.StatusServiceUnavailable {
74 t.Fatal("ls-refs was held to the pack limit")
75 }
76}
77
78// A request with a valid bearer token counts against the account, the
79// key SSH uses; anything else against the client address.
80func TestPackPrincipal(t *testing.T) {
81 s := busyServer(t)
82 alice, err := s.st.UserByUsername("alice")
83 if err != nil {
84 t.Fatal(err)
85 }
86 if err := s.st.CreateAPIToken(alice.ID, "t", store.HashToken("secret"), "read", nil, 0); err != nil {
87 t.Fatal(err)
88 }
89 r := httptest.NewRequest("POST", "/alice/app/git-upload-pack", nil)
90 r.RemoteAddr = "192.0.2.7:4000"
91 if got := s.packPrincipal(r); got != "ip:192.0.2.7" {
92 t.Fatalf("anonymous: %q", got)
93 }
94 r.Header.Set("Authorization", "Bearer wrong")
95 if got := s.packPrincipal(r); got != "ip:192.0.2.7" {
96 t.Fatalf("bad token: %q", got)
97 }
98 r6 := httptest.NewRequest("POST", "/alice/app/git-upload-pack", nil)
99 r6.RemoteAddr = "[2001:db8:1:2:3:4:5:6]:4000"
100 if got := s.packPrincipal(r6); got != "ip:2001:db8:1:2::/64" {
101 t.Fatalf("anonymous IPv6: %q", got)
102 }
103 want := "user:" + strconv.FormatInt(alice.ID, 10)
104 r.Header.Set("Authorization", "Bearer secret")
105 if got := s.packPrincipal(r); got != want {
106 t.Fatalf("token: %q, want %q", got, want)
107 }
108 if err := s.st.CreateWebSession(store.HashToken("sess"), alice.ID, time.Hour); err != nil {
109 t.Fatal(err)
110 }
111 r = httptest.NewRequest("POST", "/alice/app/git-upload-pack", nil)
112 r.RemoteAddr = "192.0.2.7:4000"
113 r.AddCookie(&http.Cookie{Name: sessionCookie, Value: "sess"})
114 if got := s.packPrincipal(r); got != want {
115 t.Fatalf("session: %q, want %q", got, want)
116 }
117}
118
119// stuckClient is a connection whose client sent the start of a request
120// body and then stopped sending. Reads block until a read deadline is
121// set; the response is discarded.
122type stuckClient struct {
123 header http.Header
124 head string // the part of the body that was sent
125 cut chan struct{}
126 once sync.Once
127}
128
129func (c *stuckClient) Header() http.Header { return c.header }
130func (c *stuckClient) WriteHeader(int) {}
131func (c *stuckClient) Write(b []byte) (int, error) { return len(b), nil }
132func (c *stuckClient) Read(b []byte) (int, error) {
133 if c.head != "" {
134 n := copy(b, c.head)
135 c.head = c.head[n:]
136 return n, nil
137 }
138 <-c.cut
139 return 0, os.ErrDeadlineExceeded
140}
141func (c *stuckClient) Close() error { return nil }
142func (c *stuckClient) SetReadDeadline(time.Time) error {
143 c.once.Do(func() { close(c.cut) })
144 return nil
145}
146
147// limitedServer is a server with a pack limit and an empty alice/app.
148func limitedServer(t *testing.T) *Server {
149 t.Helper()
150 s := busyServer(t)
151 s.packs = packlimit.New(1, 0, 0, time.Second)
152 if err := gitutil.InitBare(control.RepoDir(s.cfg.Server.Root, "alice", "app"), "main", t.TempDir()); err != nil {
153 t.Fatal(err)
154 }
155 return s
156}
157
158// seed commits files of the given sizes, random and so incompressible,
159// to alice/app's main and returns the commit.
160func seed(t *testing.T, s *Server, sizes ...int) string {
161 t.Helper()
162 work := t.TempDir()
163 git := func(args ...string) string {
164 t.Helper()
165 cmd := exec.Command("git", append([]string{"-C", work, "-c", "user.name=t", "-c", "user.email=t@t"}, args...)...)
166 out, err := cmd.CombinedOutput()
167 if err != nil {
168 t.Fatalf("git %v: %v\n%s", args, err, out)
169 }
170 return strings.TrimSpace(string(out))
171 }
172 git("init", "-q", "-b", "main")
173 for i, n := range sizes {
174 b := make([]byte, n)
175 rand.Read(b)
176 if err := os.WriteFile(filepath.Join(work, strconv.Itoa(i)), b, 0o644); err != nil {
177 t.Fatal(err)
178 }
179 }
180 git("add", ".")
181 git("commit", "-q", "-m", "seed")
182 git("push", "-q", control.RepoDir(s.cfg.Server.Root, "alice", "app"), "main")
183 return git("rev-parse", "HEAD")
184}
185
186// fetchBody is a protocol v0 request for sha's whole history.
187func fetchBody(sha string) string {
188 pkt := func(s string) string { return fmt.Sprintf("%04x%s", len(s)+4, s) }
189 return pkt("want "+sha+" side-band-64k ofs-delta\n") + "0000" + pkt("done\n")
190}
191
192// ended requires the handler to finish within limit and its slot to be
193// free.
194func ended(t *testing.T, s *Server, finished <-chan struct{}, limit time.Duration) {
195 t.Helper()
196 select {
197 case <-finished:
198 case <-time.After(limit):
199 t.Fatal("fetch still running")
200 }
201 hold, err := s.packs.Acquire(nil, "ip:elsewhere")
202 if err != nil {
203 t.Fatalf("slot not released after the kill: %v", err)
204 }
205 hold()
206}
207
208func stallAfter(t *testing.T, d time.Duration) {
209 old := packlimit.StallDeadline
210 packlimit.StallDeadline = d
211 t.Cleanup(func() { packlimit.StallDeadline = old })
212}
213
214// A client that stops sending its body is cut after StallDeadline.
215func TestFetchKilledWhenClientStopsSending(t *testing.T) {
216 stallAfter(t, 200*time.Millisecond)
217 s := limitedServer(t)
218 // Half a pkt-line: git waits for the rest.
219 c := &stuckClient{header: http.Header{}, head: "0032want 0123456789abcdef", cut: make(chan struct{})}
220 r := httptest.NewRequest("POST", "/alice/app/git-upload-pack", c)
221 r.SetPathValue("owner", "alice")
222 r.SetPathValue("repo", "app")
223 finished := make(chan struct{})
224 go func() {
225 s.uploadPack(c, r)
226 close(finished)
227 }()
228 ended(t, s, finished, 5*time.Second)
229}
230
231// A client that leaves ends git at once, even while git is busy with
232// neither its input nor its output: here a pack-objects hook that
233// sleeps, with the whole request read and the response discarded.
234func TestFetchKilledWhenClientLeaves(t *testing.T) {
235 s := limitedServer(t)
236 sha := seed(t, s, 10)
237 hook := filepath.Join(t.TempDir(), "hook")
238 if err := os.WriteFile(hook, []byte("#!/bin/sh\nsleep 60\n"), 0o755); err != nil {
239 t.Fatal(err)
240 }
241 t.Setenv("GIT_CONFIG_COUNT", "1")
242 t.Setenv("GIT_CONFIG_KEY_0", "uploadpack.packObjectsHook")
243 t.Setenv("GIT_CONFIG_VALUE_0", hook)
244 ctx, cancel := context.WithCancel(context.Background())
245 defer cancel()
246 r := httptest.NewRequestWithContext(ctx, "POST", "/alice/app/git-upload-pack", strings.NewReader(fetchBody(sha)))
247 r.SetPathValue("owner", "alice")
248 r.SetPathValue("repo", "app")
249 finished := make(chan struct{})
250 go func() {
251 s.uploadPack(httptest.NewRecorder(), r)
252 close(finished)
253 }()
254 time.Sleep(500 * time.Millisecond)
255 cancel()
256 // Unkilled, git runs until the hook's sleep ends.
257 ended(t, s, finished, 5*time.Second)
258}
259
260// A client that stops reading the response is cut after StallDeadline:
261// the write blocked on its full socket fails at the deadline set then.
262func TestFetchKilledWhenClientStopsReading(t *testing.T) {
263 stallAfter(t, 500*time.Millisecond)
264 s := limitedServer(t)
265 sizes := make([]int, 16)
266 for i := range sizes {
267 sizes[i] = 2 << 20
268 }
269 sha := seed(t, s, sizes...)
270 finished := make(chan struct{})
271 srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
272 r.SetPathValue("owner", "alice")
273 r.SetPathValue("repo", "app")
274 s.uploadPack(w, r)
275 close(finished)
276 }))
277 defer srv.Close()
278 conn, err := net.Dial("tcp", srv.Listener.Addr().String())
279 if err != nil {
280 t.Fatal(err)
281 }
282 defer conn.Close()
283 conn.(*net.TCPConn).SetReadBuffer(4096)
284 body := fetchBody(sha)
285 fmt.Fprintf(conn, "POST /alice/app/git-upload-pack HTTP/1.1\r\nHost: x\r\nContent-Length: %d\r\n\r\n%s", len(body), body)
286 // The 32MB pack cannot fit in the socket buffers; nothing is read.
287 ended(t, s, finished, 20*time.Second)
288}
289
290func getArchive(s *Server, w http.ResponseWriter, r *http.Request) {
291 r.SetPathValue("owner", "alice")
292 r.SetPathValue("repo", "app")
293 r.SetPathValue("file", "main.tar.gz")
294 s.archive(w, r)
295}
296
297// A web archive takes a pack slot: busy is 503, and the slot is free
298// again once the archive is written.
299func TestArchiveTakesASlot(t *testing.T) {
300 s := limitedServer(t)
301 seed(t, s, 10)
302 hold, err := s.packs.Acquire(nil, "ip:elsewhere")
303 if err != nil {
304 t.Fatal(err)
305 }
306 w := httptest.NewRecorder()
307 getArchive(s, w, httptest.NewRequest("GET", "/alice/app/archive/main.tar.gz", nil))
308 if w.Code != http.StatusServiceUnavailable || w.Header().Get("Retry-After") == "" {
309 t.Fatalf("busy: status %d, Retry-After %q", w.Code, w.Header().Get("Retry-After"))
310 }
311 hold()
312 w = httptest.NewRecorder()
313 getArchive(s, w, httptest.NewRequest("GET", "/alice/app/archive/main.tar.gz", nil))
314 if w.Code != http.StatusOK || w.Header().Get("Content-Type") != "application/gzip" || w.Body.Len() == 0 {
315 t.Fatalf("free: status %d, type %q, %d bytes", w.Code, w.Header().Get("Content-Type"), w.Body.Len())
316 }
317 hold, err = s.packs.Acquire(nil, "ip:elsewhere")
318 if err != nil {
319 t.Fatalf("slot not released after the archive: %v", err)
320 }
321 hold()
322}
323
324// An archive whose client stops reading is cut after StallDeadline.
325func TestArchiveKilledWhenClientStopsReading(t *testing.T) {
326 stallAfter(t, 500*time.Millisecond)
327 s := limitedServer(t)
328 sizes := make([]int, 16)
329 for i := range sizes {
330 sizes[i] = 2 << 20
331 }
332 seed(t, s, sizes...)
333 finished := make(chan struct{})
334 srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
335 getArchive(s, w, r)
336 close(finished)
337 }))
338 defer srv.Close()
339 conn, err := net.Dial("tcp", srv.Listener.Addr().String())
340 if err != nil {
341 t.Fatal(err)
342 }
343 defer conn.Close()
344 conn.(*net.TCPConn).SetReadBuffer(4096)
345 fmt.Fprintf(conn, "GET /alice/app/archive/main.tar.gz HTTP/1.1\r\nHost: x\r\n\r\n")
346 // The 32MB archive cannot fit in the socket buffers; nothing is read.
347 ended(t, s, finished, 20*time.Second)
348}
internal/httpd/routes_test.go +4 −4
@@ -14,7 +14,7 @@ import (
14func TestViewOnlyHasNoMutatingRoutes(t *testing.T) { 14func TestViewOnlyHasNoMutatingRoutes(t *testing.T) {
15 cfg := config.Default() 15 cfg := config.Default()
16 cfg.Web.Mode = "view_only" 16 cfg.Web.Mode = "view_only"
17 s := New(cfg, nil) 17 s := New(cfg, nil, nil)
18 18
19 for _, r := range s.Routes() { 19 for _, r := range s.Routes() {
20 if r.Mutating { 20 if r.Mutating {
@@ -38,7 +38,7 @@ func TestViewOnlyHasNoMutatingRoutes(t *testing.T) {
38// TestAPIRouteGating: the API route exists only when [api] enabled = true. 38// TestAPIRouteGating: the API route exists only when [api] enabled = true.
39func TestAPIRouteGating(t *testing.T) { 39func TestAPIRouteGating(t *testing.T) {
40 has := func(cfg config.Config) bool { 40 has := func(cfg config.Config) bool {
41 for _, r := range New(cfg, nil).Routes() { 41 for _, r := range New(cfg, nil, nil).Routes() {
42 if r.Pattern == "/api/v1/cmd" { 42 if r.Pattern == "/api/v1/cmd" {
43 return true 43 return true
44 } 44 }
@@ -60,7 +60,7 @@ func TestAPIRouteGating(t *testing.T) {
60func TestAccountsModeHasLoginRoute(t *testing.T) { 60func TestAccountsModeHasLoginRoute(t *testing.T) {
61 cfg := config.Default() 61 cfg := config.Default()
62 cfg.Web.Mode = "accounts" 62 cfg.Web.Mode = "accounts"
63 s := New(cfg, nil) 63 s := New(cfg, nil, nil)
64 found := false 64 found := false
65 for _, r := range s.Routes() { 65 for _, r := range s.Routes() {
66 if r.Pattern == "/login" { 66 if r.Pattern == "/login" {
@@ -82,7 +82,7 @@ func TestTopLevelRouteWordsAreReserved(t *testing.T) {
82 // (config.Default() leaves it "closed"); open it so this walk actually 82 // (config.Default() leaves it "closed"); open it so this walk actually
83 // reaches the route the production instance runs with. 83 // reaches the route the production instance runs with.
84 cfg.Registration.Mode = "open" 84 cfg.Registration.Mode = "open"
85 s := New(cfg, nil) 85 s := New(cfg, nil, nil)
86 for _, r := range s.Routes() { 86 for _, r := range s.Routes() {
87 seg := strings.TrimPrefix(r.Pattern, "/") 87 seg := strings.TrimPrefix(r.Pattern, "/")
88 seg, _, _ = strings.Cut(seg, "/") 88 seg, _, _ = strings.Cut(seg, "/")
internal/httpd/smart.go +100 −6
@@ -6,18 +6,24 @@
6package httpd 6package httpd
7 7
8import ( 8import (
9 "bufio"
9 "compress/gzip" 10 "compress/gzip"
11 "errors"
10 "fmt" 12 "fmt"
11 "io" 13 "io"
12 "net" 14 "net"
13 "net/http" 15 "net/http"
14 "os" 16 "os"
15 "os/exec" 17 "os/exec"
18 "strconv"
16 "strings" 19 "strings"
17 "sync" 20 "sync"
21 "time"
18 22
19 "gitbay.org/gitbay/internal/config" 23 "gitbay.org/gitbay/internal/config"
20 "gitbay.org/gitbay/internal/control" 24 "gitbay.org/gitbay/internal/control"
25 "gitbay.org/gitbay/internal/gitutil"
26 "gitbay.org/gitbay/internal/packlimit"
21 "gitbay.org/gitbay/internal/store" 27 "gitbay.org/gitbay/internal/store"
22 "gitbay.org/gitbay/internal/toolpath" 28 "gitbay.org/gitbay/internal/toolpath"
23) 29)
@@ -25,15 +31,16 @@ import (
25type Server struct { 31type Server struct {
26 cfg config.Config 32 cfg config.Config
27 st *store.Store 33 st *store.Store
34 packs *packlimit.Limiter
28 apiLimit *apiLimiter 35 apiLimit *apiLimiter
29 proxies []*net.IPNet // http.trusted_proxies, parsed once 36 proxies []*net.IPNet // http.trusted_proxies, parsed once
30 stopping chan struct{} // closed by Stop 37 stopping chan struct{} // closed by Stop
31 stopOnce sync.Once 38 stopOnce sync.Once
32} 39}
33 40
34func New(cfg config.Config, st *store.Store) *Server { 41func New(cfg config.Config, st *store.Store, packs *packlimit.Limiter) *Server {
35 proxies, _ := cfg.HTTP.TrustedProxyNets() // validated at config load 42 proxies, _ := cfg.HTTP.TrustedProxyNets() // validated at config load
36 return &Server{cfg: cfg, st: st, apiLimit: newAPILimiter(cfg.Limits.APIRate), proxies: proxies, 43 return &Server{cfg: cfg, st: st, packs: packs, apiLimit: newAPILimiter(cfg.Limits.APIRate), proxies: proxies,
37 stopping: make(chan struct{})} 44 stopping: make(chan struct{})}
38} 45}
39 46
@@ -135,14 +142,101 @@ func (s *Server) uploadPack(w http.ResponseWriter, r *http.Request) {
135 defer gz.Close() 142 defer gz.Close()
136 body = gz 143 body = gz
137 } 144 }
145 br := bufio.NewReader(body)
146 cancel := r.Context().Done()
147 out := io.Writer(w)
148 if !lsRefs(br) {
149 o, kill, finish, ok := s.packSlot(w, r)
150 if !ok {
151 return
152 }
153 // Deferred before git runs, so it fires after git has exited
154 // and been waited for.
155 defer finish()
156 out, cancel = o, kill
157 }
138 w.Header().Set("Content-Type", "application/x-git-upload-pack-result") 158 w.Header().Set("Content-Type", "application/x-git-upload-pack-result")
139 w.Header().Set("Cache-Control", "no-cache") 159 w.Header().Set("Cache-Control", "no-cache")
140 dir := control.RepoDir(s.cfg.Server.Root, repo.OwnerName, repo.Name) 160 dir := control.RepoDir(s.cfg.Server.Root, repo.OwnerName, repo.Name)
141 cmd := exec.CommandContext(r.Context(), toolpath.Look("git"), "upload-pack", "--stateless-rpc", dir) 161 cmd := exec.Command(toolpath.Look("git"), "-c", "uploadpack.keepAlive=5", "upload-pack", "--stateless-rpc", dir)
142 cmd.Env = append(os.Environ(), gitProtocolEnv(r)...) 162 cmd.Env = append(os.Environ(), gitProtocolEnv(r)...)
143 cmd.Stdin = body 163 cmd.Stdin = br
144 cmd.Stdout = w 164 cmd.Stdout = out
145 cmd.Run() 165 gitutil.RunUntil(cmd, cancel)
166}
167
168// packSlot takes a pack-generation slot for r, answering 503 with
169// Retry-After when none comes free. A queued request waits at most the
170// limiter's wait, and Stop ends the wait so it does not hold up a
171// restart's drain; net/http notices a departed client only after the
172// body is read, so that rarely ends it. On success out is w watched for
173// stalls, kill closes when git must stop — the client left, or no write
174// to it completed for packlimit.StallDeadline — and finish, called once
175// git has exited, releases the slot.
176func (s *Server) packSlot(w http.ResponseWriter, r *http.Request) (out io.Writer, kill <-chan struct{}, finish func(), ok bool) {
177 principal := s.packPrincipal(r)
178 release, err := s.packs.Acquire(s.until(r), principal)
179 if err != nil {
180 s.packs.Refused("http", principal, err)
181 msg := "the server is restarting; try again in a minute"
182 if errors.Is(err, packlimit.ErrBusy) {
183 msg = "the server is busy: it is at its limit of concurrent clones and fetches; try again in a minute"
184 }
185 w.Header().Set("Retry-After", "30")
186 http.Error(w, msg, http.StatusServiceUnavailable)
187 return nil, nil, nil, false
188 }
189 out, stalled, unwatch := s.packs.Watch(w)
190 killed := make(chan struct{})
191 finished := make(chan struct{})
192 exited := make(chan struct{})
193 go func() {
194 defer close(exited)
195 select {
196 case <-finished:
197 return
198 case <-r.Context().Done():
199 case <-stalled:
200 // A write blocked on a client that stopped reading, or a
201 // read of a body it stopped sending, outlives git;
202 // expired deadlines end both copies, so Wait returns.
203 rc := http.NewResponseController(w)
204 rc.SetReadDeadline(time.Now())
205 rc.SetWriteDeadline(time.Now())
206 }
207 close(killed)
208 }()
209 return out, killed, func() {
210 // The watcher must not touch w once the handler has returned.
211 close(finished)
212 <-exited
213 unwatch()
214 release()
215 }, true
216}
217
218// lsRefs reports whether a protocol v2 request is a ref listing, which
219// generates no pack. Its first pkt-line is "command=ls-refs".
220func lsRefs(br *bufio.Reader) bool {
221 const want = "command=ls-refs"
222 head, err := br.Peek(4 + len(want))
223 return err == nil && string(head[4:]) == want
224}
225
226// packPrincipal is who a fetch is counted against: the account when the
227// request carries a valid bearer token or web session, the same key SSH
228// uses, so switching transport buys no extra slots; otherwise the
229// client address, an IPv6 one by its /64.
230func (s *Server) packPrincipal(r *http.Request) string {
231 if tok, ok := strings.CutPrefix(r.Header.Get("Authorization"), "Bearer "); ok && strings.TrimSpace(tok) != "" {
232 if u, _, err := s.st.APITokenUser(store.HashToken(strings.TrimSpace(tok))); err == nil {
233 return "user:" + strconv.FormatInt(u.ID, 10)
234 }
235 }
236 if u := s.viewer(r); u.ID != 0 {
237 return "user:" + strconv.FormatInt(u.ID, 10)
238 }
239 return packlimit.AddrPrincipal(s.clientIP(r))
146} 240}
147 241
148// gitProtocolEnv forwards the client's protocol negotiation header so 242// gitProtocolEnv forwards the client's protocol negotiation header so
internal/httpd/web.go +6 −1
@@ -2340,10 +2340,15 @@ func (s *Server) archive(w http.ResponseWriter, r *http.Request) {
2340 s.notFound(w, r) 2340 s.notFound(w, r)
2341 return 2341 return
2342 } 2342 }
2343 out, kill, finish, ok := s.packSlot(w, r)
2344 if !ok {
2345 return
2346 }
2347 defer finish()
2343 prefix := fmt.Sprintf("%s-%s", p.Repo.Name, ref) 2348 prefix := fmt.Sprintf("%s-%s", p.Repo.Name, ref)
2344 w.Header().Set("Content-Type", "application/gzip") 2349 w.Header().Set("Content-Type", "application/gzip")
2345 w.Header().Set("Content-Disposition", fmt.Sprintf("attachment; filename=%q", prefix+".tar.gz")) 2350 w.Header().Set("Content-Disposition", fmt.Sprintf("attachment; filename=%q", prefix+".tar.gz"))
2346 gitutil.Archive(p.Dir, ref, prefix, w) 2351 gitutil.ArchiveUntil(p.Dir, ref, prefix, out, kill)
2347} 2352}
2348 2353
2349func policyCanAdmin(u store.User, repo store.Repo, grant string) bool { 2354func policyCanAdmin(u store.User, repo store.Repo, grant string) bool {
internal/packlimit/packlimit.go added +205
@@ -0,0 +1,205 @@
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.
6// 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.
8package packlimit
9
10import (
11 "errors"
12 "log/slog"
13 "net/netip"
14 "strings"
15 "sync"
16 "time"
17)
18
19var (
20 ErrBusy = errors.New("the server is busy generating packs for other clients; try again in a minute")
21 ErrGone = errors.New("client went away while queued")
22)
23
24type Limiter struct {
25 max, per, queue int
26 wait time.Duration
27
28 // Principals starting with class may hold at most classCap slots
29 // between them; classCap 0 is no class cap.
30 class string
31 classCap int
32
33 mu sync.Mutex
34 running int
35 classHeld int
36 queued int
37 held map[string]int // running, per principal
38 waiting map[string]int // queued, per principal
39 changed chan struct{} // closed and replaced on every release
40 warned map[string]time.Time // last refusal logged, per transport
41}
42
43// New returns a limiter, or nil — no limit — when max is not positive.
44func New(max, per, queue int, wait time.Duration) *Limiter {
45 if max <= 0 {
46 return nil
47 }
48 return &Limiter{max: max, per: per, queue: queue, wait: wait,
49 held: map[string]int{}, waiting: map[string]int{}, changed: make(chan struct{}),
50 warned: map[string]time.Time{}}
51}
52
53// Refused logs that a request on transport was turned away with err, at
54// most once a minute per transport. It names the principal's class
55// (user or ip), never the principal: an address is personal data.
56func (l *Limiter) Refused(transport, principal string, err error) {
57 if l == nil {
58 return
59 }
60 now := time.Now()
61 l.mu.Lock()
62 last, seen := l.warned[transport]
63 if seen && now.Sub(last) < time.Minute {
64 l.mu.Unlock()
65 return
66 }
67 l.warned[transport] = now
68 l.mu.Unlock()
69 class, _, _ := strings.Cut(principal, ":")
70 reason := "busy"
71 if errors.Is(err, ErrGone) {
72 reason = "gone"
73 }
74 slog.Warn("pack limit: request turned away (logged at most once a minute per transport)",
75 "transport", transport, "class", class, "reason", reason)
76}
77
78// CapClass caps the slots that principals starting with prefix may hold
79// between them. Call it before the limiter is in use.
80func (l *Limiter) CapClass(prefix string, n int) {
81 if l == nil {
82 return
83 }
84 l.class, l.classCap = prefix, n
85}
86
87// AddrPrincipal is the principal for an unauthenticated client at addr:
88// an IPv4 address as is, an IPv6 address by its /64, since one host
89// commonly holds a whole /64. An address that does not parse is used
90// as given.
91func AddrPrincipal(addr string) string {
92 a, err := netip.ParseAddr(addr)
93 if err != nil {
94 return "ip:" + addr
95 }
96 a = a.WithZone("").Unmap()
97 if a.Is4() {
98 return "ip:" + a.String()
99 }
100 return "ip:" + netip.PrefixFrom(a, 64).Masked().String()
101}
102
103// Acquire takes a slot for principal, queueing when none is free.
104// done, when it closes, ends the wait. Once Acquire returns a nil
105// error, the caller holds the slot and must call release — once git
106// has exited — regardless of what its own context has done since:
107// done closing after that point does not release the slot on the
108// caller's behalf.
109func (l *Limiter) Acquire(done <-chan struct{}, principal string) (release func(), err error) {
110 if l == nil {
111 return func() {}, nil
112 }
113 l.mu.Lock()
114 if l.fits(principal) {
115 l.take(principal)
116 l.mu.Unlock()
117 return l.releaser(principal), nil
118 }
119 if l.queued >= l.queue || (l.per > 0 && l.waiting[principal] >= l.per) {
120 l.mu.Unlock()
121 return nil, ErrBusy
122 }
123 l.queued++
124 l.waiting[principal]++
125 l.mu.Unlock()
126 defer func() {
127 l.mu.Lock()
128 l.queued--
129 if l.waiting[principal]--; l.waiting[principal] == 0 {
130 delete(l.waiting, principal)
131 }
132 l.mu.Unlock()
133 }()
134
135 timer := time.NewTimer(l.wait)
136 defer timer.Stop()
137 for {
138 l.mu.Lock()
139 // changed and done can both be ready at once — a slot can
140 // free up at the same moment the caller gives up. select
141 // among the wake sources would then pick between them at
142 // random, so re-check done first, under the lock, on every
143 // pass: this makes the limiter prefer ErrGone whenever both
144 // are ready, instead of leaving it to chance which one a
145 // given pass observes.
146 select {
147 case <-done:
148 l.mu.Unlock()
149 return nil, ErrGone
150 default:
151 }
152 if l.fits(principal) {
153 l.take(principal)
154 l.mu.Unlock()
155 return l.releaser(principal), nil
156 }
157 changed := l.changed
158 l.mu.Unlock()
159 select {
160 case <-changed:
161 case <-timer.C:
162 return nil, ErrBusy
163 case <-done:
164 return nil, ErrGone
165 }
166 }
167}
168
169func (l *Limiter) fits(principal string) bool {
170 if l.inClass(principal) && l.classHeld >= l.classCap {
171 return false
172 }
173 return l.running < l.max && (l.per <= 0 || l.held[principal] < l.per)
174}
175
176func (l *Limiter) inClass(principal string) bool {
177 return l.classCap > 0 && strings.HasPrefix(principal, l.class)
178}
179
180func (l *Limiter) take(principal string) {
181 l.running++
182 l.held[principal]++
183 if l.inClass(principal) {
184 l.classHeld++
185 }
186}
187
188func (l *Limiter) releaser(principal string) func() {
189 var once sync.Once
190 return func() {
191 once.Do(func() {
192 l.mu.Lock()
193 defer l.mu.Unlock()
194 l.running--
195 if l.inClass(principal) {
196 l.classHeld--
197 }
198 if l.held[principal]--; l.held[principal] == 0 {
199 delete(l.held, principal)
200 }
201 close(l.changed)
202 l.changed = make(chan struct{})
203 })
204 }
205}
internal/packlimit/packlimit_test.go added +315
@@ -0,0 +1,315 @@
1package packlimit
2
3import (
4 "bytes"
5 "errors"
6 "log/slog"
7 "math"
8 "strings"
9 "testing"
10 "time"
11)
12
13func TestGlobalCap(t *testing.T) {
14 l := New(2, 0, 0, time.Second)
15 r1, err1 := l.Acquire(nil, "a")
16 r2, err2 := l.Acquire(nil, "b")
17 if err1 != nil || err2 != nil {
18 t.Fatal(err1, err2)
19 }
20 if _, err := l.Acquire(nil, "c"); !errors.Is(err, ErrBusy) {
21 t.Fatalf("third with no queue: %v", err)
22 }
23 r1()
24 r1() // a second release is a no-op
25 r3, err := l.Acquire(nil, "c")
26 if err != nil {
27 t.Fatal(err)
28 }
29 if _, err := l.Acquire(nil, "d"); !errors.Is(err, ErrBusy) {
30 t.Fatalf("double release freed two slots: %v", err)
31 }
32 r2()
33 r3()
34}
35
36func TestPerPrincipalCap(t *testing.T) {
37 l := New(4, 1, 4, 50*time.Millisecond)
38 ra, err := l.Acquire(nil, "a")
39 if err != nil {
40 t.Fatal(err)
41 }
42 defer ra()
43 if _, err := l.Acquire(nil, "a"); !errors.Is(err, ErrBusy) {
44 t.Fatalf("second for a: %v", err)
45 }
46 rb, err := l.Acquire(nil, "b")
47 if err != nil {
48 t.Fatalf("b blocked by a: %v", err)
49 }
50 rb()
51}
52
53func TestWaiterGetsReleasedSlot(t *testing.T) {
54 l := New(1, 0, 1, 5*time.Second)
55 r1, _ := l.Acquire(nil, "a")
56 got := make(chan error, 1)
57 go func() {
58 r, err := l.Acquire(nil, "b")
59 if err == nil {
60 r()
61 }
62 got <- err
63 }()
64 waitQueued(t, l, 1)
65 r1()
66 select {
67 case err := <-got:
68 if err != nil {
69 t.Fatal(err)
70 }
71 case <-time.After(2 * time.Second):
72 t.Fatal("waiter never got the slot")
73 }
74}
75
76func TestQueueIsBounded(t *testing.T) {
77 l := New(1, 0, 1, 5*time.Second)
78 r1, _ := l.Acquire(nil, "a")
79 defer r1()
80 go l.Acquire(nil, "b")
81 waitQueued(t, l, 1)
82 if _, err := l.Acquire(nil, "c"); !errors.Is(err, ErrBusy) {
83 t.Fatalf("queue over its bound: %v", err)
84 }
85}
86
87// A principal cannot fill the queue on its own.
88func TestPrincipalQueueIsBounded(t *testing.T) {
89 l := New(1, 1, 8, 5*time.Second)
90 r1, _ := l.Acquire(nil, "x")
91 defer r1()
92 go l.Acquire(nil, "a")
93 waitQueued(t, l, 1)
94 if _, err := l.Acquire(nil, "a"); !errors.Is(err, ErrBusy) {
95 t.Fatalf("second waiter for a: %v", err)
96 }
97}
98
99func TestClientGoneWhileQueued(t *testing.T) {
100 l := New(1, 0, 1, 5*time.Second)
101 r1, _ := l.Acquire(nil, "a")
102 done := make(chan struct{})
103 close(done)
104 if _, err := l.Acquire(done, "b"); !errors.Is(err, ErrGone) {
105 t.Fatalf("got %v, want ErrGone", err)
106 }
107 assertQueueEmpty(t, l)
108 r1()
109 assertHeldEmpty(t, l)
110}
111
112func TestWaitRunsOut(t *testing.T) {
113 l := New(1, 0, 1, 20*time.Millisecond)
114 r1, _ := l.Acquire(nil, "a")
115 if _, err := l.Acquire(nil, "b"); !errors.Is(err, ErrBusy) {
116 t.Fatalf("got %v, want ErrBusy", err)
117 }
118 assertQueueEmpty(t, l)
119 r1()
120 assertHeldEmpty(t, l)
121}
122
123func assertQueueEmpty(t *testing.T, l *Limiter) {
124 t.Helper()
125 l.mu.Lock()
126 defer l.mu.Unlock()
127 if l.queued != 0 || len(l.waiting) != 0 {
128 t.Fatalf("queue not cleaned up: queued=%d waiting=%v", l.queued, l.waiting)
129 }
130}
131
132func assertHeldEmpty(t *testing.T, l *Limiter) {
133 t.Helper()
134 l.mu.Lock()
135 defer l.mu.Unlock()
136 if len(l.held) != 0 {
137 t.Fatalf("held not cleaned up: %v", l.held)
138 }
139}
140
141func TestNilLimiterNeverWaits(t *testing.T) {
142 var l *Limiter
143 if l = New(0, 1, 1, time.Second); l != nil {
144 t.Fatal("max 0 should mean no limit")
145 }
146 r, err := l.Acquire(nil, "a")
147 if err != nil {
148 t.Fatal(err)
149 }
150 r()
151}
152
153// A waiter whose done channel closes just as a slot frees up must not be
154// granted the slot: it has to see ErrGone, and the slot must go to
155// someone else instead of leaking to an abandoned caller.
156func TestGivenUpWaiterNeverGetsSlot(t *testing.T) {
157 l := New(1, 0, 1, 10*time.Second)
158 _, _ = l.Acquire(nil, "a") // holds the only slot
159
160 done := make(chan struct{})
161 got := make(chan error, 1)
162 go func() {
163 r, err := l.Acquire(done, "b")
164 if err == nil {
165 r()
166 }
167 got <- err
168 }()
169 waitQueued(t, l, 1)
170
171 // Close done and free a's slot in the same critical section, so
172 // changed and done become ready to b's select at the same instant
173 // — the exact race the done-check-before-fits ordering in Acquire
174 // has to win, whichever the select picks.
175 l.mu.Lock()
176 close(done)
177 l.running--
178 delete(l.held, "a")
179 close(l.changed)
180 l.changed = make(chan struct{})
181 l.mu.Unlock()
182
183 select {
184 case err := <-got:
185 if !errors.Is(err, ErrGone) {
186 t.Fatalf("got %v, want ErrGone", err)
187 }
188 case <-time.After(2 * time.Second):
189 t.Fatal("b never returned")
190 }
191
192 // The slot must still be free for someone else: b must not hold it.
193 r2, err := l.Acquire(nil, "c")
194 if err != nil {
195 t.Fatalf("slot leaked to the abandoned waiter: %v", err)
196 }
197 r2()
198}
199
200func waitQueued(t *testing.T, l *Limiter, n int) {
201 t.Helper()
202 deadline := time.Now().Add(2 * time.Second)
203 for time.Now().Before(deadline) {
204 l.mu.Lock()
205 q := l.queued
206 l.mu.Unlock()
207 if q == n {
208 return
209 }
210 time.Sleep(time.Millisecond)
211 }
212 t.Fatalf("queue never reached %d", n)
213}
214
215// config maps pack_queue = -1 to math.MaxInt: waiters are not turned
216// away for want of queue room.
217func TestUnboundedQueue(t *testing.T) {
218 l := New(1, 0, math.MaxInt, 5*time.Second)
219 r1, _ := l.Acquire(nil, "a")
220 done := make(chan struct{})
221 defer close(done)
222 for i := 0; i < 64; i++ {
223 go l.Acquire(done, "b")
224 }
225 waitQueued(t, l, 64)
226 r1()
227}
228
229// Anonymous clients cannot hold every slot: with max 3 and "ip:" capped
230// at 2, a third anonymous request queues and an account still gets in.
231func TestClassCap(t *testing.T) {
232 l := New(3, 0, 4, 5*time.Second)
233 l.CapClass("ip:", 2)
234 r1, err1 := l.Acquire(nil, "ip:192.0.2.1")
235 r2, err2 := l.Acquire(nil, "ip:192.0.2.2")
236 if err1 != nil || err2 != nil {
237 t.Fatal(err1, err2)
238 }
239 got := make(chan error, 1)
240 go func() {
241 r, err := l.Acquire(nil, "ip:192.0.2.3")
242 if err == nil {
243 r()
244 }
245 got <- err
246 }()
247 waitQueued(t, l, 1)
248 ru, err := l.Acquire(nil, "user:1")
249 if err != nil {
250 t.Fatalf("account refused the free slot: %v", err)
251 }
252 select {
253 case err := <-got:
254 t.Fatalf("third anonymous request did not queue: %v", err)
255 default:
256 }
257 r1()
258 select {
259 case err := <-got:
260 if err != nil {
261 t.Fatal(err)
262 }
263 case <-time.After(2 * time.Second):
264 t.Fatal("queued anonymous request never got the freed slot")
265 }
266 r2()
267 ru()
268 assertHeldEmpty(t, l)
269 if l.classHeld != 0 {
270 t.Fatalf("classHeld = %d after every release", l.classHeld)
271 }
272}
273
274func TestAddrPrincipal(t *testing.T) {
275 for in, want := range map[string]string{
276 "192.0.2.7": "ip:192.0.2.7",
277 "::ffff:192.0.2.7": "ip:192.0.2.7",
278 "2001:db8:1:2:3:4:5:6": "ip:2001:db8:1:2::/64",
279 "2001:db8:1:2:ffff::1": "ip:2001:db8:1:2::/64",
280 "fe80::1%en0": "ip:fe80::/64",
281 "not-an-address": "ip:not-an-address",
282 } {
283 if got := AddrPrincipal(in); got != want {
284 t.Errorf("AddrPrincipal(%q) = %q, want %q", in, got, want)
285 }
286 }
287}
288
289// A refusal is logged once a minute per transport, with the principal's
290// class and never its address.
291func TestRefusedLogsOncePerTransport(t *testing.T) {
292 var buf bytes.Buffer
293 old := slog.Default()
294 slog.SetDefault(slog.New(slog.NewTextHandler(&buf, nil)))
295 t.Cleanup(func() { slog.SetDefault(old) })
296
297 l := New(1, 0, 0, time.Second)
298 l.Refused("http", "ip:192.0.2.7", ErrBusy)
299 l.Refused("http", "ip:192.0.2.8", ErrBusy)
300 l.Refused("ssh", "user:4", ErrGone)
301 out := buf.String()
302 if n := strings.Count(out, "\n"); n != 2 {
303 t.Fatalf("%d lines, want 2:\n%s", n, out)
304 }
305 for _, want := range []string{"transport=http class=ip reason=busy", "transport=ssh class=user reason=gone"} {
306 if !strings.Contains(out, want) {
307 t.Errorf("missing %q in:\n%s", want, out)
308 }
309 }
310 if strings.Contains(out, "192.0.2") || strings.Contains(out, "user:4") {
311 t.Fatalf("principal logged:\n%s", out)
312 }
313 var none *Limiter
314 none.Refused("git", "ip:x", ErrBusy)
315}
internal/packlimit/watch.go added +61
@@ -0,0 +1,61 @@
1package packlimit
2
3import (
4 "io"
5 "sync"
6 "sync/atomic"
7 "time"
8)
9
10// StallDeadline is how long a limited transport may go without
11// completing a write to its client before it is killed. upload-pack
12// sends a keepalive every five seconds while it prepares a pack.
13var StallDeadline = 2 * time.Minute
14
15// Watch wraps w, a transport's writer to its client, so that a client
16// that stops reading does not hold its slot for as long as its
17// connection stays open. stalled closes once no write to w has
18// completed for StallDeadline; stop ends the watch. With no limit in
19// force (a nil Limiter) nothing is watched: w comes back as is and
20// stalled never closes.
21func (l *Limiter) Watch(w io.Writer) (out io.Writer, stalled <-chan struct{}, stop func()) {
22 if l == nil {
23 return w, nil, func() {}
24 }
25 deadline := StallDeadline
26 pw := &progressWriter{w: w}
27 pw.last.Store(time.Now().UnixNano())
28 st := make(chan struct{})
29 quit := make(chan struct{})
30 go func() {
31 t := time.NewTicker(deadline / 4)
32 defer t.Stop()
33 for {
34 select {
35 case <-quit:
36 return
37 case <-t.C:
38 if time.Since(time.Unix(0, pw.last.Load())) >= deadline {
39 close(st)
40 return
41 }
42 }
43 }
44 }()
45 var once sync.Once
46 return pw, st, func() { once.Do(func() { close(quit) }) }
47}
48
49// progressWriter records when a write to the client last completed.
50type progressWriter struct {
51 w io.Writer
52 last atomic.Int64 // unix nanoseconds
53}
54
55func (p *progressWriter) Write(b []byte) (int, error) {
56 n, err := p.w.Write(b)
57 if n > 0 {
58 p.last.Store(time.Now().UnixNano())
59 }
60 return n, err
61}
internal/packlimit/watch_test.go added +40
@@ -0,0 +1,40 @@
1package packlimit
2
3import (
4 "io"
5 "testing"
6 "time"
7)
8
9func TestWatchStalls(t *testing.T) {
10 old := StallDeadline
11 StallDeadline = 100 * time.Millisecond
12 t.Cleanup(func() { StallDeadline = old })
13 l := New(1, 0, 0, time.Second)
14 w, stalled, stop := l.Watch(io.Discard)
15 defer stop()
16 // Writes keep it alive past the deadline.
17 for range 6 {
18 time.Sleep(40 * time.Millisecond)
19 w.Write([]byte("x"))
20 select {
21 case <-stalled:
22 t.Fatal("stalled while writing")
23 default:
24 }
25 }
26 select {
27 case <-stalled:
28 case <-time.After(2 * time.Second):
29 t.Fatal("no stall after writes stopped")
30 }
31}
32
33func TestWatchWithoutLimit(t *testing.T) {
34 var l *Limiter
35 w, stalled, stop := l.Watch(io.Discard)
36 defer stop()
37 if w != io.Discard || stalled != nil {
38 t.Fatal("a nil limiter watched")
39 }
40}
internal/sshd/refusal_test.go +174 −1
@@ -2,12 +2,17 @@ package sshd
2 2
3import ( 3import (
4 "bytes" 4 "bytes"
5 "io"
6 "os"
5 "path/filepath" 7 "path/filepath"
6 "strings" 8 "strings"
7 "testing" 9 "testing"
10 "time"
8 11
9 "gitbay.org/gitbay/internal/config" 12 "gitbay.org/gitbay/internal/config"
10 "gitbay.org/gitbay/internal/control" 13 "gitbay.org/gitbay/internal/control"
14 "gitbay.org/gitbay/internal/gitutil"
15 "gitbay.org/gitbay/internal/packlimit"
11 "gitbay.org/gitbay/internal/protocol" 16 "gitbay.org/gitbay/internal/protocol"
12 "gitbay.org/gitbay/internal/store" 17 "gitbay.org/gitbay/internal/store"
13) 18)
@@ -51,7 +56,7 @@ func TestRefusedPushIsAudited(t *testing.T) {
51 bob.Pending = pending 56 bob.Pending = pending
52 key := store.SSHKey{Scope: "full", Fingerprint: "SHA256:test"} 57 key := store.SSHKey{Scope: "full", Fingerprint: "SHA256:test"}
53 var out, errOut bytes.Buffer 58 var out, errOut bytes.Buffer
54 code := Exec(cfg, st, bob, key, control.Term{}, "git-receive-pack alice/app", 59 code := Exec(cfg, st, nil, bob, key, control.Term{}, "git-receive-pack alice/app",
55 strings.NewReader(""), &out, &errOut, nil, nil, nil) 60 strings.NewReader(""), &out, &errOut, nil, nil, nil)
56 if code != protocol.ExitDenied { 61 if code != protocol.ExitDenied {
57 t.Fatalf("pending %v: exit %d: %s", pending, code, errOut.String()) 62 t.Fatalf("pending %v: exit %d: %s", pending, code, errOut.String())
@@ -66,3 +71,171 @@ func TestRefusedPushIsAudited(t *testing.T) {
66 } 71 }
67 } 72 }
68} 73}
74
75func TestCloneRefusedWhenPackSlotsAreFull(t *testing.T) {
76 cfg, st, bob := execFixture(t)
77 packs := packlimit.New(1, 0, 0, time.Second)
78 hold, err := packs.Acquire(nil, "ip:elsewhere")
79 if err != nil {
80 t.Fatal(err)
81 }
82 defer hold()
83 key := store.SSHKey{Scope: "full", Fingerprint: "SHA256:test"}
84 for _, service := range []string{"git-upload-pack", "git-upload-archive"} {
85 var out, errOut bytes.Buffer
86 code := Exec(cfg, st, packs, bob, key, control.Term{}, service+" alice/app",
87 strings.NewReader(""), &out, &errOut, nil, nil, nil)
88 if code != protocol.ExitFailure || !strings.Contains(errOut.String(), "busy") {
89 t.Fatalf("%s: exit %d: %q", service, code, errOut.String())
90 }
91 }
92}
93
94// A push takes no pack slot: it runs while every slot is held. A clone
95// gives its slot back once git has exited.
96func TestPushBypassesPackLimitAndCloneReleasesSlot(t *testing.T) {
97 cfg, st, _ := execFixture(t)
98 alice, err := st.UserByUsername("alice")
99 if err != nil {
100 t.Fatal(err)
101 }
102 if err := gitutil.InitBare(control.RepoDir(cfg.Server.Root, "alice", "app"), "main", t.TempDir()); err != nil {
103 t.Fatal(err)
104 }
105 key := store.SSHKey{Scope: "full", Fingerprint: "SHA256:test"}
106 packs := packlimit.New(1, 0, 0, time.Second)
107
108 var out, errOut bytes.Buffer
109 if code := Exec(cfg, st, packs, alice, key, control.Term{}, "git-upload-pack alice/app",
110 strings.NewReader("0000"), &out, &errOut, nil, nil, nil); code != protocol.ExitOK {
111 t.Fatalf("clone: exit %d: %s", code, errOut.String())
112 }
113 hold, err := packs.Acquire(nil, "ip:elsewhere")
114 if err != nil {
115 t.Fatalf("slot not released after the clone: %v", err)
116 }
117 defer hold()
118
119 out.Reset()
120 errOut.Reset()
121 if code := Exec(cfg, st, packs, alice, key, control.Term{}, "git-receive-pack alice/app",
122 strings.NewReader("0000"), &out, &errOut, nil, nil, nil); code != protocol.ExitOK {
123 t.Fatalf("push with slots full: exit %d: %s", code, errOut.String())
124 }
125}
126
127// cloneFixture adds an empty bare alice/app on disk and returns alice.
128func cloneFixture(t *testing.T) (config.Config, *store.Store, store.User) {
129 t.Helper()
130 cfg, st, _ := execFixture(t)
131 alice, err := st.UserByUsername("alice")
132 if err != nil {
133 t.Fatal(err)
134 }
135 if err := gitutil.InitBare(control.RepoDir(cfg.Server.Root, "alice", "app"), "main", t.TempDir()); err != nil {
136 t.Fatal(err)
137 }
138 return cfg, st, alice
139}
140
141// silentStdin is a client that sends nothing and never hangs up. It is
142// an *os.File, so git reads it directly: git exits only when killed.
143func silentStdin(t *testing.T) *os.File {
144 t.Helper()
145 r, w, err := os.Pipe()
146 if err != nil {
147 t.Fatal(err)
148 }
149 t.Cleanup(func() { r.Close(); w.Close() })
150 return r
151}
152
153// killedClone runs a clone of alice/app with the given channels and
154// requires it to end, killed, within five seconds, with its slot free.
155func killedClone(t *testing.T, stdout io.Writer, done, stopping, revoked <-chan struct{}) {
156 t.Helper()
157 cfg, st, alice := cloneFixture(t)
158 key := store.SSHKey{Scope: "full", Fingerprint: "SHA256:test"}
159 packs := packlimit.New(1, 0, 0, time.Second)
160 codec := make(chan int, 1)
161 go func() {
162 codec <- Exec(cfg, st, packs, alice, key, control.Term{}, "git-upload-pack alice/app",
163 silentStdin(t), stdout, io.Discard, done, stopping, revoked)
164 }()
165 select {
166 case code := <-codec:
167 if code != protocol.ExitFailure {
168 t.Fatalf("exit %d, want the clone killed", code)
169 }
170 case <-time.After(5 * time.Second):
171 t.Fatal("clone still running")
172 }
173 hold, err := packs.Acquire(nil, "ip:elsewhere")
174 if err != nil {
175 t.Fatalf("slot not released after the kill: %v", err)
176 }
177 hold()
178}
179
180func closed() <-chan struct{} {
181 c := make(chan struct{})
182 close(c)
183 return c
184}
185
186func TestCloneKilledWhenClientLeaves(t *testing.T) {
187 killedClone(t, io.Discard, closed(), nil, nil)
188}
189
190// A revoked key ends a clone even during a restart.
191func TestCloneKilledWhenKeyRevoked(t *testing.T) {
192 killedClone(t, io.Discard, nil, closed(), closed())
193}
194
195// A client that stops reading is cut after packlimit.StallDeadline.
196func TestCloneKilledWhenClientStopsReading(t *testing.T) {
197 old := packlimit.StallDeadline
198 packlimit.StallDeadline = 200 * time.Millisecond
199 t.Cleanup(func() { packlimit.StallDeadline = old })
200 r, w := io.Pipe()
201 t.Cleanup(func() { r.Close() })
202 killedClone(t, w, nil, nil, nil)
203}
204
205// On a restart (done and stopping both closed) a running clone finishes.
206func TestCloneRunsOnDuringRestart(t *testing.T) {
207 cfg, st, alice := cloneFixture(t)
208 key := store.SSHKey{Scope: "full", Fingerprint: "SHA256:test"}
209 packs := packlimit.New(1, 0, 0, time.Second)
210 var errOut bytes.Buffer
211 if code := Exec(cfg, st, packs, alice, key, control.Term{}, "git-upload-pack alice/app",
212 strings.NewReader("0000"), io.Discard, &errOut, closed(), closed(), nil); code != protocol.ExitOK {
213 t.Fatalf("exit %d: %s", code, errOut.String())
214 }
215}
216
217// A request the key may not make is refused before it reaches the
218// limiter: not found, never busy.
219func TestRefusedCloneStaysOffLimiter(t *testing.T) {
220 cfg, st, bob := execFixture(t)
221 alice, err := st.UserByUsername("alice")
222 if err != nil {
223 t.Fatal(err)
224 }
225 if _, err := st.CreateRepo("user", alice.ID, "secret", "private"); err != nil {
226 t.Fatal(err)
227 }
228 packs := packlimit.New(1, 0, 0, time.Second)
229 hold, err := packs.Acquire(nil, "ip:elsewhere")
230 if err != nil {
231 t.Fatal(err)
232 }
233 defer hold()
234 key := store.SSHKey{Scope: "full", Fingerprint: "SHA256:test"}
235 var out, errOut bytes.Buffer
236 code := Exec(cfg, st, packs, bob, key, control.Term{}, "git-upload-pack alice/secret",
237 strings.NewReader(""), &out, &errOut, nil, nil, nil)
238 if code != protocol.ExitNotFound || strings.Contains(errOut.String(), "busy") {
239 t.Fatalf("exit %d: %q", code, errOut.String())
240 }
241}
internal/sshd/sshd.go +75 −8
@@ -29,6 +29,7 @@ import (
29 "gitbay.org/gitbay/internal/control" 29 "gitbay.org/gitbay/internal/control"
30 "gitbay.org/gitbay/internal/gitutil" 30 "gitbay.org/gitbay/internal/gitutil"
31 "gitbay.org/gitbay/internal/hookd" 31 "gitbay.org/gitbay/internal/hookd"
32 "gitbay.org/gitbay/internal/packlimit"
32 "gitbay.org/gitbay/internal/policy" 33 "gitbay.org/gitbay/internal/policy"
33 "gitbay.org/gitbay/internal/protocol" 34 "gitbay.org/gitbay/internal/protocol"
34 "gitbay.org/gitbay/internal/store" 35 "gitbay.org/gitbay/internal/store"
@@ -37,6 +38,7 @@ import (
37type Server struct { 38type Server struct {
38 cfg config.Config 39 cfg config.Config
39 st *store.Store 40 st *store.Store
41 packs *packlimit.Limiter
40 sshCfg *ssh.ServerConfig 42 sshCfg *ssh.ServerConfig
41 authLimiter *rateLimiter 43 authLimiter *rateLimiter
42 sessions sync.WaitGroup // accepted connections still being served 44 sessions sync.WaitGroup // accepted connections still being served
@@ -68,8 +70,8 @@ func (c *conn) cut() {
68 c.net.Close() 70 c.net.Close()
69} 71}
70 72
71func New(cfg config.Config, st *store.Store) (*Server, error) { 73func New(cfg config.Config, st *store.Store, packs *packlimit.Limiter) (*Server, error) {
72 s := &Server{cfg: cfg, st: st, authLimiter: newRateLimiter(cfg.Limits.SSHAuthRate, time.Minute), conns: map[*conn]struct{}{}, stopping: make(chan struct{})} 74 s := &Server{cfg: cfg, st: st, packs: packs, authLimiter: newRateLimiter(cfg.Limits.SSHAuthRate, time.Minute), conns: map[*conn]struct{}{}, stopping: make(chan struct{})}
73 75
74 sc := &ssh.ServerConfig{ 76 sc := &ssh.ServerConfig{
75 PublicKeyCallback: s.authenticate, 77 PublicKeyCallback: s.authenticate,
@@ -423,7 +425,7 @@ func (s *Server) runExec(c *conn, sconn *ssh.ServerConn, ch ssh.Channel, term co
423 return protocol.ExitDenied 425 return protocol.ExitDenied
424 } 426 }
425 _ = s.st.TouchSSHKey(keyID) 427 _ = s.st.TouchSSHKey(keyID)
426 return Exec(s.cfg, s.st, user, key, term, cmdline, ch, ch, ch.Stderr(), done, s.stopping, c.revoked) 428 return Exec(s.cfg, s.st, s.packs, user, key, term, cmdline, ch, ch, ch.Stderr(), done, s.stopping, c.revoked)
427} 429}
428 430
429// runAnonymous handles a session from an unregistered key: the register 431// runAnonymous handles a session from an unregistered key: the register
@@ -459,7 +461,7 @@ func (s *Server) runAnonymous(ch ssh.Channel, keyB64, cmdline string) int {
459// Exec runs one SSH exec command line for an authenticated key. It is the 461// Exec runs one SSH exec command line for an authenticated key. It is the
460// single dispatch path shared by the embedded listener and the system-sshd 462// single dispatch path shared by the embedded listener and the system-sshd
461// forced command (gitbayd shell). Closing revoked kills a git transport. 463// forced command (gitbayd shell). Closing revoked kills a git transport.
462func Exec(cfg config.Config, st *store.Store, user store.User, key store.SSHKey, term control.Term, cmdline string, 464func Exec(cfg config.Config, st *store.Store, packs *packlimit.Limiter, user store.User, key store.SSHKey, term control.Term, cmdline string,
463 stdin io.Reader, stdout, stderr io.Writer, done, stopping, revoked <-chan struct{}) int { 465 stdin io.Reader, stdout, stderr io.Writer, done, stopping, revoked <-chan struct{}) int {
464 if user.Disabled { 466 if user.Disabled {
465 fmt.Fprintln(stderr, "this account is disabled; contact the instance admin") 467 fmt.Fprintln(stderr, "this account is disabled; contact the instance admin")
@@ -477,7 +479,7 @@ func Exec(cfg config.Config, st *store.Store, user store.User, key store.SSHKey,
477 if user.Pending { 479 if user.Pending {
478 fmt.Fprintln(stderr, "your account is not active yet: verify your email first") 480 fmt.Fprintln(stderr, "your account is not active yet: verify your email first")
479 } else { 481 } else {
480 code = runGit(cfg, st, user, key.Scope, argv, stdin, stdout, stderr, revoked) 482 code = runGit(cfg, st, packs, user, key.Scope, argv, stdin, stdout, stderr, done, stopping, revoked)
481 } 483 }
482 // A refused push is a refused write, audited like one. runGit 484 // A refused push is a refused write, audited like one. runGit
483 // refuses only with the path as the one argument, so argv[1:] 485 // refuses only with the path as the one argument, so argv[1:]
@@ -515,8 +517,8 @@ func Exec(cfg config.Config, st *store.Store, user store.User, key store.SSHKey,
515} 517}
516 518
517// runGit streams a git transport service after access checks. 519// runGit streams a git transport service after access checks.
518func runGit(cfg config.Config, st *store.Store, user store.User, scope string, argv []string, 520func runGit(cfg config.Config, st *store.Store, packs *packlimit.Limiter, user store.User, scope string, argv []string,
519 stdin io.Reader, stdout, stderr io.Writer, revoked <-chan struct{}) int { 521 stdin io.Reader, stdout, stderr io.Writer, done, stopping, revoked <-chan struct{}) int {
520 service := argv[0] 522 service := argv[0]
521 if len(argv) != 2 { 523 if len(argv) != 2 {
522 fmt.Fprintf(stderr, "usage: %s <path>\n", service) 524 fmt.Fprintf(stderr, "usage: %s <path>\n", service)
@@ -601,7 +603,72 @@ func runGit(cfg config.Config, st *store.Store, user store.User, scope string, a
601 defer st.DeletePushToken(token) 603 defer st.DeletePushToken(token)
602 env = append(env, hookd.EnvToken+"="+token) 604 env = append(env, hookd.EnvToken+"="+token)
603 } 605 }
604 if err := gitutil.Transport(service, dir, stdin, stdout, stderr, env, maxPack, revoked); err != nil { 606 cancel := revoked
607 if !write {
608 // 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
624 }
625 // Deferred before Transport runs, so it fires after git has
626 // exited and been waited for.
627 defer release()
628 // A client that stops reading would hold its slot for as long
629 // as its channel stays open.
630 client := stdout
631 var stalled <-chan struct{}
632 var unwatch func()
633 stdout, stalled, unwatch = packs.Watch(client)
634 defer unwatch()
635 kill := make(chan struct{})
636 finished := make(chan struct{})
637 defer close(finished)
638 go func() {
639 left := done
640 for {
641 select {
642 case <-finished:
643 return
644 case <-revoked:
645 case <-left:
646 select {
647 case <-stopping:
648 // done closes on a restart too; a clone already
649 // running finishes then. Only a departed client
650 // ends it.
651 left = nil
652 continue
653 default:
654 }
655 case <-stalled:
656 close(kill)
657 // A write blocked on the client's window outlives
658 // git; closing the channel ends it and the stdin copy,
659 // so Transport's Wait returns.
660 if c, ok := client.(io.Closer); ok {
661 c.Close()
662 }
663 return
664 }
665 close(kill)
666 return
667 }
668 }()
669 cancel = kill
670 }
671 if err := gitutil.Transport(service, dir, stdin, stdout, stderr, env, maxPack, cancel); err != nil {
605 return protocol.ExitFailure 672 return protocol.ExitFailure
606 } 673 }
607 return protocol.ExitOK 674 return protocol.ExitOK
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) 67 srv, err := New(cfg, st, 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) 248 srv, err := New(cfg, st, nil)
249 if err != nil { 249 if err != nil {
250 t.Fatal(err) 250 t.Fatal(err)
251 } 251 }