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=.
208208 Organizations are not capped.
209209- =max_bytes_per_user= (0, unlimited) — disk the account's own
210210 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.
212241
213242** [git_daemon]
214243- =enabled= (false), =port= (9418) — the anonymous =git://= listener.
@@ -258,8 +287,12 @@ meaningful — an unverified address never produces a =verified= badge.
258287The audit log is the security feed (events are the product feed): every
259288successful mutating command with its argv and source credential (SSH key
260289fingerprint or API), every refused one (exit 3 or 4) as =refused
261<command>=, refused pushes as =refused git-receive-pack=, registrations,
262admin actions, force-pushes, and auth failures/throttling. A refusal row
290<command>=, refused pushes as =refused git-receive-pack= (the access
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
263296keeps the flag names and the first positional, not the values.
264297Refusals are recorded up to ten a minute per account and 600 a minute
265298across 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
272305verify= opens the store as other admin commands do, applying pending
273306migrations, so run it with the binary that matches the daemon. It
274307recomputes 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
276312rows is not a break. Rows written before the chain existed are counted
277313and skipped; when every row is such a row, verify warns and exits 1,
278314since 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
282318afterwards under the freed ids. The database cannot show either. The
283319daemon logs every row it writes to its journal, outside the database
284320(=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
286323host =gitbayd admin= commands, and by =gitbayd shell= when =ssh.mode =
287324"system"=, are not copied to the journal.
288325
.gitbay/wiki/Architecture/09-Controls.org +2 −2
@@ -69,7 +69,7 @@ chapter names of OWASP ASVS 4.0 where one fits.
6969| Security-relevant writes audited | in place | every successful mutating command (=control.go=) |
7070| Authentication failures audited | in place | =auth.failed=, =auth.throttled= |
7171| 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 |
7373
7474** Communications and integrations (V9, V10, V12)
7575
@@ -96,7 +96,7 @@ chapter names of OWASP ASVS 4.0 where one fits.
9696| Control | Status | Evidence |
9797|---------------------------------------------+----------+------------------------------------------------------------------|
9898| 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 |
100100| Service hardening | in place | systemd sandboxing ([[file:03-Deployment.org][3]]) |
101101| Backups offsite and append-only | in place | restic with append-only credentials (documented) |
102102| 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.
1313| #259 | Recovery | No restore has been exercised; the drill is written (Admin wiki) and not yet run | high |
1414| #260 | CI network | Builds share the runner's source address; no egress policy | medium |
1515| #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 |
1816| #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 |
1918
2019* Not filed
2120
2221| Area | Gap | Severity |
2322|-------+-------------------------------------------------------------------------------------------------------------+----------|
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 |
2528
2629* Questions an auditor will ask that have no answer yet
2730
.gitbay/wiki/Performance.org +17 −2
@@ -45,5 +45,20 @@ this scale never appears in a profile.
4545
4646The practical ceiling on this hardware is concurrent pack generation:
4747full 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
49cores, not with changes to gitbay.
48clone ran git at ~156% CPU). =limits.pack_concurrency= bounds how many
49run at once across SSH, HTTP and git://, with a queue behind it (see
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:
249249 escape what is there. #144 covers the missing isolation.
250250- *Timing and traffic analysis.* Token comparison is a hash index lookup
251251 by design, but nothing has been measured.
252- *Denial of service by resource exhaustion* beyond rate: large pushes,
253 pathological diffs, deep histories, zip bombs in LFS.
252- *Denial of service by resource exhaustion* beyond rate. Concurrent
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.
254256
255257A sweep is a point in time. This section says what a reader should not
256258assume has been checked.
@@ -271,13 +273,15 @@ assume has been checked.
271273 present as objects no archived ref names, and a repository's refs may
272274 be newer than the database snapshot (see [[Admin]]).
273275- 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
275 or removed row show as a break under =gitbayd admin audit verify=,
276 except at the end: removing the newest rows, and writing new rows
277 under their freed ids, leaves a valid chain. Only comparing verify's
278 last id and hash with the daemon's journal copy shows it, and rows
279 written outside the daemon (=gitbayd shell= under =ssh.mode =
280 "system"=, host =gitbayd admin= commands) have no journal copy.
276 the daemon user's access can change it. The hash chain is unkeyed:
277 whoever can write the database can edit a row and recompute every
278 later hash. =gitbayd admin audit verify= catches an edited or removed
279 row only when the later hashes were not recomputed, and never catches
280 removing the newest rows or writing new rows under their freed ids.
281 Comparing verify's last id and hash with the daemon's journal copy is
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.
281285- A global signature-verification epoch over-invalidates the cache on any
282286 trust-input change. Correct, not a leak; a performance tradeoff.
283287- 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
9999 before it (migration 0064), and the daemon logs a copy of every row it
100100 writes to its journal. =gitbayd admin audit verify= prints the row
101101 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
103 shows only by comparing that last id and hash with the journal (#275).
102 was edited or whose predecessor was removed. The chain is unkeyed:
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).
104107- Refused mutating commands (exit 3 or 4) are audited as =refused
105108 <command>=, and refused pushes as =refused git-receive-pack=, keeping
106109 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
108111 =refused.throttled= row stands for the rest of the minute. Under
109112 =ssh.mode = "system"= each =gitbayd shell= connection counts
110113 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).
111118- Audit retention deletes by id, up to the newest row older than the
112119 retention, so a clock step back cannot leave a gap in the chain (#275).
113120- =dashboard= and =feed= print activity as sentences
@@ -217,6 +224,29 @@ missing, =gitbayd admin backup --verify <archive>= names it, and
217224 repository on the host. (#259)
218225- =gitbayd admin secrets init= and =rotate= hold an flock on =<key
219226 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).
220250
221251* v1.36.0 — 2026-09-23
222252
cmd/gitbayd/main.go +13 −3
@@ -32,6 +32,7 @@ import (
3232 "gitbay.org/gitbay/internal/httpd"
3333 "gitbay.org/gitbay/internal/mirror"
3434 "gitbay.org/gitbay/internal/notify"
35 "gitbay.org/gitbay/internal/packlimit"
3536 "gitbay.org/gitbay/internal/push"
3637 "gitbay.org/gitbay/internal/seal"
3738 "gitbay.org/gitbay/internal/sshd"
@@ -228,11 +229,20 @@ func serveCmd() *cobra.Command {
228229 return control.RepoDir(cfg.Server.Root, owner, name)
229230 }, buildinfo.String()).Run(whCtx)
230231
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
231241 errCh := make(chan error, 3)
232242 var sshSrv *sshd.Server
233243 var sshLn, gitLn net.Listener
234244 if cfg.SSH.Mode == "embedded" {
235 srv, err := sshd.New(cfg, st)
245 srv, err := sshd.New(cfg, st, packs)
236246 if err != nil {
237247 return err
238248 }
@@ -249,7 +259,7 @@ func serveCmd() *cobra.Command {
249259 slog.Info("ssh handled by host sshd (ssh.mode = system)")
250260 }
251261
252 web := httpd.New(cfg, st)
262 web := httpd.New(cfg, st, packs)
253263 // Header and idle timeouts bound what an idle or slow client can
254264 // hold open. No write timeout: archives and upload-pack stream
255265 // for as long as they take (#104).
@@ -338,7 +348,7 @@ func serveCmd() *cobra.Command {
338348 }
339349 slog.Info("git-daemon listening", "addr", gln.Addr())
340350 gitLn = gln
341 go func() { errCh <- gitd.New(cfg, st).Serve(gln) }()
351 go func() { errCh <- gitd.New(cfg, st, packs).Serve(gln) }()
342352 }
343353
344354 select {
cmd/gitbayd/system.go +3 −1
@@ -99,7 +99,9 @@ func shellCmd() *cobra.Command {
9999 fmt.Fprintf(os.Stderr, "gitbay control plane: interactive shells are not available.\nTry: ssh <host> help\n")
100100 os.Exit(protocol.ExitUsage)
101101 }
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)
103105 st.Close()
104106 os.Exit(code)
105107 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 (
77 "encoding/pem"
88 "errors"
99 "fmt"
10 "math"
1011 "net"
1112 "os"
1213 "path/filepath"
@@ -24,6 +25,15 @@ import (
2425// and webhook deliveries.
2526const DefaultWriteRate = 60
2627
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
2737type Config struct {
2838 Server Server `toml:"server"`
2939 SSH SSH `toml:"ssh"`
@@ -214,6 +224,41 @@ type Limits struct {
214224 // account.
215225 MaxReposPerUser int `toml:"max_repos_per_user"`
216226 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
217262}
218263
219264type Mail struct {
@@ -444,6 +489,11 @@ func (c Config) Validate() error {
444489 if c.Limits.MaxReposPerUser < 0 || c.Limits.MaxBytesPerUser < 0 || c.Limits.MaxSnippetsPerUser < 0 {
445490 errs = append(errs, errors.New("limits.max_repos_per_user, max_bytes_per_user and max_snippets_per_user must not be negative"))
446491 }
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 }
447497 if c.Push.Enabled {
448498 for _, f := range []struct{ name, val string }{
449499 {"push.key_file", c.Push.KeyFile},
internal/config/config_test.go +22
@@ -6,12 +6,14 @@ import (
66 "crypto/rand"
77 "crypto/x509"
88 "encoding/pem"
9 "math"
910 "os"
1011 "path/filepath"
1112 "strings"
1213 "testing"
1314
1415 "filippo.io/age"
16 "time"
1517)
1618
1719func writeConfig(t *testing.T, body string) string {
@@ -46,12 +48,32 @@ func TestLoadMinimal(t *testing.T) {
4648 }
4749}
4850
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
4966func TestContradictions(t *testing.T) {
5067 cases := []struct {
5168 name string
5269 body string
5370 wantErr string
5471 }{
72 {
73 "bad pack_queue_wait",
74 minimal + "\n[limits]\npack_queue_wait = \"soon\"\n",
75 "limits.pack_queue_wait",
76 },
5577 {
5678 "registration open without smtp",
5779 minimal + "\n[registration]\nmode = \"open\"\n",
internal/gitd/gitd.go +43 −12
@@ -7,24 +7,26 @@ import (
77 "fmt"
88 "io"
99 "net"
10 "os"
11 "os/exec"
1210 "strconv"
1311 "strings"
1412 "time"
1513
1614 "gitbay.org/gitbay/internal/config"
1715 "gitbay.org/gitbay/internal/control"
16 "gitbay.org/gitbay/internal/gitutil"
17 "gitbay.org/gitbay/internal/packlimit"
1818 "gitbay.org/gitbay/internal/store"
19 "gitbay.org/gitbay/internal/toolpath"
2019)
2120
2221type Server struct {
23 cfg config.Config
24 st *store.Store
22 cfg config.Config
23 st *store.Store
24 packs *packlimit.Limiter
2525}
2626
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}
2830
2931func (s *Server) Serve(ln net.Listener) error {
3032 for {
@@ -68,13 +70,42 @@ func (s *Server) handle(conn net.Conn) {
6870 return
6971 }
7072
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
71101 dir := control.RepoDir(s.cfg.Server.Root, repo.OwnerName, repo.Name)
72 cmd := exec.Command(toolpath.Look("git"), "upload-pack", dir)
73 cmd.Env = append(os.Environ(), protoEnv...)
74 cmd.Stdin = conn
75 cmd.Stdout = conn
76 cmd.Stderr = io.Discard
77 cmd.Run()
102 gitutil.Transport("git-upload-pack", dir, conn, out, io.Discard, protoEnv, 0, kill)
103}
104
105// principal is the pack-limit principal for a client at addr.
106func principal(addr net.Addr) string {
107 host, _, _ := net.SplitHostPort(addr.String())
108 return packlimit.AddrPrincipal(host)
78109}
79110
80111func 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
4545 if service == "git-receive-pack" && maxPack > 0 {
4646 args = []string{"-c", fmt.Sprintf("receive.maxInputSize=%d", maxPack)}
4747 }
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 }
4853 args = append(args, strings.TrimPrefix(service, "git-"), repoPath)
4954 default:
5055 return fmt.Errorf("unknown service %q", service)
@@ -54,6 +59,13 @@ func Transport(service, repoPath string, stdin io.Reader, stdout, errW io.Writer
5459 cmd.Stdin = stdin
5560 cmd.Stdout = stdout
5661 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 {
5769 ownProcessGroup(cmd)
5870 if err := cmd.Start(); err != nil {
5971 return err
internal/gitutil/read.go +6 −1
@@ -132,12 +132,17 @@ var ErrArchiveTooLarge = errors.New("archive exceeds the size limit")
132132// Archive streams a tar.gz of ref to w, within archiveTimeout and
133133// MaxArchiveBytes. Past either, git is killed and the error says which.
134134func 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 {
135140 ctx, cancel := context.WithTimeout(context.Background(), archiveTimeout)
136141 defer cancel()
137142 cmd := exec.CommandContext(ctx, toolpath.Look("git"), "-C", dir, "archive", "--format=tar.gz", "--prefix="+prefix+"/", "--end-of-options", ref)
138143 lw := &cappedWriter{w: w, left: MaxArchiveBytes, stop: cancel}
139144 cmd.Stdout = lw
140 err := cmd.Run()
145 err := RunUntil(cmd, stop)
141146 switch {
142147 case lw.exceeded:
143148 return ErrArchiveTooLarge
internal/hookd/hookd.go +38 −10
@@ -122,8 +122,10 @@ func (s *Server) handle(conn net.Conn) {
122122 defer conn.Close()
123123 dec := json.NewDecoder(conn)
124124 enc := json.NewEncoder(conn)
125 if err := checkPeer(conn); err != nil {
125 if err := peerCheck(conn); err != nil {
126126 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()})
127129 enc.Encode(Response{Allow: false, Message: "hook socket: " + err.Error()})
128130 return
129131 }
@@ -132,7 +134,9 @@ func (s *Server) handle(conn net.Conn) {
132134 enc.Encode(Response{Allow: false, Message: "bad hook request"})
133135 return
134136 }
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})
136140 enc.Encode(Response{Allow: false, Message: msg})
137141 return
138142 }
@@ -149,18 +153,42 @@ func (s *Server) handle(conn net.Conn) {
149153
150154// authorize ties a request to a receive-pack sshd started: its token
151155// 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) {
153159 if req.Token == "" {
154 return "push not started by this server"
160 return 0, "push not started by this server"
155161 }
156162 tok, err := s.st.PushTokenByHash(store.HashToken(req.Token))
157163 if err != nil {
158 return "push not started by this server"
164 return 0, "push not started by this server"
159165 }
160166 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"
162168 }
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})
164192}
165193
166194func (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) {
170198 return
171199 }
172200 if msg := policy.CheckPush(repo, req.Updates); msg != "" {
173 enc.Encode(Response{Allow: false, Message: msg})
201 s.refusePush(enc, req, repo, msg)
174202 return
175203 }
176204 if msg := s.releaseAnchors(repo, req.Updates); msg != "" {
177 enc.Encode(Response{Allow: false, Message: msg})
205 s.refusePush(enc, req, repo, msg)
178206 return
179207 }
180208 if !repo.Settings.RequireSignedCommits {
@@ -220,7 +248,7 @@ func (s *Server) preReceive(req Request, dec *json.Decoder, enc *json.Encoder) {
220248 }
221249 }
222250 if refusal != "" {
223 enc.Encode(Response{Allow: false, Message: refusal})
251 s.refusePush(enc, req, repo, refusal)
224252 return
225253 }
226254 enc.Encode(Response{Allow: true})
internal/hookd/socket_test.go +92
@@ -1,12 +1,17 @@
11package hookd
22
33import (
4 "encoding/json"
5 "errors"
6 "fmt"
7 "net"
48 "os"
59 "path/filepath"
610 "strings"
711 "testing"
812
913 "gitbay.org/gitbay/internal/config"
14 "gitbay.org/gitbay/internal/policy"
1015 "gitbay.org/gitbay/internal/store"
1116)
1217
@@ -103,3 +108,90 @@ func TestHookRequestNeedsItsPushToken(t *testing.T) {
103108 t.Fatalf("finished push: %+v, %v", resp, err)
104109 }
105110}
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) {
3434 t.Fatal(err)
3535 }
3636
37 s := New(config.Default(), st)
37 s := New(config.Default(), st, nil)
3838 rr := httptest.NewRecorder()
3939 req := httptest.NewRequest("GET", "/settings", nil)
4040 s.accountPage(rr, req, store.User{ID: uid, Username: "alice"})
@@ -78,7 +78,7 @@ func TestAccountSubmitNotifyPush(t *testing.T) {
7878 t.Fatal(err)
7979 }
8080 u := store.User{ID: uid, Username: "alice"}
81 s := New(config.Default(), st)
81 s := New(config.Default(), st, nil)
8282
8383 rr := submitAccountForm(t, s, u, url.Values{"field": {"notify-push"}, "push": {"on"}})
8484 if rr.Code != http.StatusSeeOther {
@@ -117,7 +117,7 @@ func TestAccountSubmitDeviceRemove(t *testing.T) {
117117 if err != nil {
118118 t.Fatal(err)
119119 }
120 s := New(config.Default(), st)
120 s := New(config.Default(), st, nil)
121121
122122 idStr := strconv.FormatInt(id, 10)
123123
@@ -158,7 +158,7 @@ func TestAccountPageMasksAShortDeviceToken(t *testing.T) {
158158 t.Fatal(err)
159159 }
160160
161 s := New(config.Default(), st)
161 s := New(config.Default(), st, nil)
162162 rr := httptest.NewRecorder()
163163 s.accountPage(rr, httptest.NewRequest("GET", "/settings", nil), store.User{ID: uid, Username: "alice"})
164164
@@ -212,7 +212,7 @@ func TestPinToggleDispatchesRepoPin(t *testing.T) {
212212 t.Fatal(err)
213213 }
214214
215 s := New(config.Default(), st)
215 s := New(config.Default(), st, nil)
216216 req := httptest.NewRequest("POST", "/alice/app/pin", nil)
217217 req.SetPathValue("owner", "alice")
218218 req.SetPathValue("repo", "app")
@@ -261,7 +261,7 @@ func TestWatchToggleCyclesThroughMuted(t *testing.T) {
261261 t.Fatal(err)
262262 }
263263
264 s := New(config.Default(), st)
264 s := New(config.Default(), st, nil)
265265 req := httptest.NewRequest("POST", "/alice/app/watch", nil)
266266 req.SetPathValue("owner", "alice")
267267 req.SetPathValue("repo", "app")
internal/httpd/admin_test.go +1 −1
@@ -48,7 +48,7 @@ func TestAdminPageShowsThePushQueue(t *testing.T) {
4848 t.Fatal(err)
4949 }
5050
51 s := New(config.Default(), st)
51 s := New(config.Default(), st, nil)
5252 rr := httptest.NewRecorder()
5353 s.adminPage(rr, httptest.NewRequest("GET", "/admin", nil), store.User{ID: uid, Username: "root", IsAdmin: true})
5454 if rr.Code != http.StatusOK {
internal/httpd/anchors_test.go +2 −2
@@ -26,7 +26,7 @@ func TestMarkdownHeadingAnchors(t *testing.T) {
2626// The stylesheet carries an ETag and a cache lifetime; a revalidation
2727// with the same tag is a 304 with no body (#132).
2828func TestStylesheetRevalidates(t *testing.T) {
29 s := New(config.Default(), nil)
29 s := New(config.Default(), nil, nil)
3030 first := httptest.NewRecorder()
3131 s.stylesheet(first, httptest.NewRequest("GET", "/static/style.css", nil))
3232 tag := first.Header().Get("ETag")
@@ -69,7 +69,7 @@ func TestStylesheetURLCarriesTheBuildHash(t *testing.T) {
6969 t.Errorf("the page does not link %s", want)
7070 }
7171
72 s := New(config.Default(), nil)
72 s := New(config.Default(), nil, nil)
7373 versioned := httptest.NewRecorder()
7474 s.stylesheet(versioned, httptest.NewRequest("GET", "/static/style.css?v="+stylesheetHash, nil))
7575 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) {
4040 // there is no second code path in routes.go for looping over both to
4141 // reach; "open" alone matches production and is enough.
4242 cfg.Registration.Mode = "open"
43 s := New(cfg, nil)
43 s := New(cfg, nil, nil)
4444
4545 for _, r := range s.Routes() {
4646 if !r.Mutating {
internal/httpd/clientip_test.go +1 −1
@@ -28,7 +28,7 @@ func TestClientIPBehindProxy(t *testing.T) {
2828 for _, tc := range cases {
2929 cfg := config.Default()
3030 cfg.HTTP.TrustedProxies = tc.proxies
31 s := New(cfg, nil)
31 s := New(cfg, nil, nil)
3232 r := httptest.NewRequest("GET", "/api/v1/read", nil)
3333 r.RemoteAddr = tc.remote
3434 if tc.xff != "" {
internal/httpd/fonts_test.go +2 −2
@@ -17,7 +17,7 @@ import (
1717// and the @font-face URLs were once maintained by hand and drifted, so
1818// gitbay.org served no web font at all (#102).
1919func TestStylesheetFontsAreServed(t *testing.T) {
20 s := New(config.Default(), nil)
20 s := New(config.Default(), nil, nil)
2121 byPattern := map[string]http.HandlerFunc{}
2222 for _, r := range s.Routes() {
2323 if r.Method == "GET" {
@@ -49,7 +49,7 @@ func TestStylesheetFontsAreServed(t *testing.T) {
4949// TestLandingImagesAreServed: every file under static/img has a route
5050// that answers 200 with an image or video type, and a video answers Range.
5151func TestLandingImagesAreServed(t *testing.T) {
52 s := New(config.Default(), nil)
52 s := New(config.Default(), nil, nil)
5353 byPattern := map[string]http.HandlerFunc{}
5454 for _, r := range s.Routes() {
5555 if r.Method == "GET" {
internal/httpd/issuecreate_test.go +5 −5
@@ -79,7 +79,7 @@ func TestIssueCreateFormHasMilestoneAndAssigneeForWriter(t *testing.T) {
7979
8080 cfg := config.Default()
8181 cfg.Web.Mode = "accounts"
82 s := New(cfg, st)
82 s := New(cfg, st, nil)
8383 req := httptest.NewRequest("GET", "/alice/app/issues/new", nil)
8484 req.SetPathValue("owner", "alice")
8585 req.SetPathValue("repo", "app")
@@ -125,7 +125,7 @@ func TestIssueCreateFormHidesMilestoneAndAssigneeForReader(t *testing.T) {
125125
126126 cfg := config.Default()
127127 cfg.Web.Mode = "accounts"
128 s := New(cfg, st)
128 s := New(cfg, st, nil)
129129 req := httptest.NewRequest("GET", "/alice/app/issues/new", nil)
130130 req.SetPathValue("owner", "alice")
131131 req.SetPathValue("repo", "app")
@@ -198,7 +198,7 @@ func TestIssueCreateSubmitReaderLabelIsDropped(t *testing.T) {
198198 t.Fatal(err)
199199 }
200200
201 s := New(config.Default(), st)
201 s := New(config.Default(), st, nil)
202202 form := url.Values{
203203 "title": {"a bug"},
204204 "body": {"steps"},
@@ -254,7 +254,7 @@ func TestIssueCreateSubmitSetsMilestoneAndAssignee(t *testing.T) {
254254 t.Fatal(err)
255255 }
256256
257 s := New(config.Default(), st)
257 s := New(config.Default(), st, nil)
258258 form := url.Values{
259259 "title": {"needs a fix"},
260260 "body": {"details"},
@@ -308,7 +308,7 @@ func TestIssueCreateSubmitBadAssigneeCreatesNothing(t *testing.T) {
308308 t.Fatal(err)
309309 }
310310
311 s := New(config.Default(), st)
311 s := New(config.Default(), st, nil)
312312 form := url.Values{
313313 "title": {"needs a fix"},
314314 "body": {"details"},
internal/httpd/logincookie_test.go +1 −1
@@ -52,7 +52,7 @@ func TestLoginNoStoreHeader(t *testing.T) {
5252 if err := st.MigrateUp(); err != nil {
5353 t.Fatal(err)
5454 }
55 s := New(config.Default(), st)
55 s := New(config.Default(), st, nil)
5656 rr := httptest.NewRecorder()
5757 req := httptest.NewRequest("GET", "/login?token=bogus", nil)
5858 s.login(rr, req)
internal/httpd/logindisabled_test.go +1 −1
@@ -40,7 +40,7 @@ func TestLoginRefusesTokenForDisabledAccount(t *testing.T) {
4040 t.Fatal(err)
4141 }
4242
43 s := New(config.Default(), st)
43 s := New(config.Default(), st, nil)
4444 rr := httptest.NewRecorder()
4545 req := httptest.NewRequest("GET", "/login?token="+tok, nil)
4646 s.login(rr, req)
internal/httpd/mrrangediff_test.go +3 −3
@@ -36,7 +36,7 @@ func TestMRRangeDiffPageRendersCommandOutput(t *testing.T) {
3636 t.Fatal(err)
3737 }
3838
39 s := New(config.Default(), st)
39 s := New(config.Default(), st, nil)
4040 req := httptest.NewRequest("GET", "/alice/app/mrs/1/range-diff", nil)
4141 req.SetPathValue("owner", "alice")
4242 req.SetPathValue("repo", "app")
@@ -180,7 +180,7 @@ func loginCookie(t *testing.T, st *store.Store, userID int64) *http.Cookie {
180180// user with no access (#269).
181181func TestMRRangeDiffPagePrivateRepo(t *testing.T) {
182182 st, cfg, alice, bob, repo, n, _, _, title := rangeDiffFixture(t)
183 s := New(cfg, st)
183 s := New(cfg, st, nil)
184184
185185 newReq := func(cookie *http.Cookie) (*httptest.ResponseRecorder, *http.Request) {
186186 req := httptest.NewRequest("GET", "/alice/secret/mrs/"+strconv.FormatInt(n, 10)+"/range-diff", nil)
@@ -232,7 +232,7 @@ func TestMRRangeDiffPagePrivateRepo(t *testing.T) {
232232// range-diff between exactly those two.
233233func TestMRRangeDiffPageFromToQuery(t *testing.T) {
234234 st, cfg, alice, _, repo, n, v1, v2, _ := rangeDiffFixture(t)
235 s := New(cfg, st)
235 s := New(cfg, st, nil)
236236 cookie := loginCookie(t, st, alice.ID)
237237
238238 newReq := func(query string) (*httptest.ResponseRecorder, *http.Request) {
internal/httpd/mrslist_test.go +1 −1
@@ -37,7 +37,7 @@ func TestMRsListContributionHintByAccess(t *testing.T) {
3737
3838 cfg := config.Default()
3939 cfg.Web.Mode = "accounts"
40 s := New(cfg, st)
40 s := New(cfg, st, nil)
4141
4242 // mrs reads the viewer through s.viewer(r), which resolves a
4343 // 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 (
1414func TestViewOnlyHasNoMutatingRoutes(t *testing.T) {
1515 cfg := config.Default()
1616 cfg.Web.Mode = "view_only"
17 s := New(cfg, nil)
17 s := New(cfg, nil, nil)
1818
1919 for _, r := range s.Routes() {
2020 if r.Mutating {
@@ -38,7 +38,7 @@ func TestViewOnlyHasNoMutatingRoutes(t *testing.T) {
3838// TestAPIRouteGating: the API route exists only when [api] enabled = true.
3939func TestAPIRouteGating(t *testing.T) {
4040 has := func(cfg config.Config) bool {
41 for _, r := range New(cfg, nil).Routes() {
41 for _, r := range New(cfg, nil, nil).Routes() {
4242 if r.Pattern == "/api/v1/cmd" {
4343 return true
4444 }
@@ -60,7 +60,7 @@ func TestAPIRouteGating(t *testing.T) {
6060func TestAccountsModeHasLoginRoute(t *testing.T) {
6161 cfg := config.Default()
6262 cfg.Web.Mode = "accounts"
63 s := New(cfg, nil)
63 s := New(cfg, nil, nil)
6464 found := false
6565 for _, r := range s.Routes() {
6666 if r.Pattern == "/login" {
@@ -82,7 +82,7 @@ func TestTopLevelRouteWordsAreReserved(t *testing.T) {
8282 // (config.Default() leaves it "closed"); open it so this walk actually
8383 // reaches the route the production instance runs with.
8484 cfg.Registration.Mode = "open"
85 s := New(cfg, nil)
85 s := New(cfg, nil, nil)
8686 for _, r := range s.Routes() {
8787 seg := strings.TrimPrefix(r.Pattern, "/")
8888 seg, _, _ = strings.Cut(seg, "/")
internal/httpd/smart.go +100 −6
@@ -6,18 +6,24 @@
66package httpd
77
88import (
9 "bufio"
910 "compress/gzip"
11 "errors"
1012 "fmt"
1113 "io"
1214 "net"
1315 "net/http"
1416 "os"
1517 "os/exec"
18 "strconv"
1619 "strings"
1720 "sync"
21 "time"
1822
1923 "gitbay.org/gitbay/internal/config"
2024 "gitbay.org/gitbay/internal/control"
25 "gitbay.org/gitbay/internal/gitutil"
26 "gitbay.org/gitbay/internal/packlimit"
2127 "gitbay.org/gitbay/internal/store"
2228 "gitbay.org/gitbay/internal/toolpath"
2329)
@@ -25,15 +31,16 @@ import (
2531type Server struct {
2632 cfg config.Config
2733 st *store.Store
34 packs *packlimit.Limiter
2835 apiLimit *apiLimiter
2936 proxies []*net.IPNet // http.trusted_proxies, parsed once
3037 stopping chan struct{} // closed by Stop
3138 stopOnce sync.Once
3239}
3340
34func New(cfg config.Config, st *store.Store) *Server {
41func New(cfg config.Config, st *store.Store, packs *packlimit.Limiter) *Server {
3542 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,
3744 stopping: make(chan struct{})}
3845}
3946
@@ -135,14 +142,101 @@ func (s *Server) uploadPack(w http.ResponseWriter, r *http.Request) {
135142 defer gz.Close()
136143 body = gz
137144 }
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 }
138158 w.Header().Set("Content-Type", "application/x-git-upload-pack-result")
139159 w.Header().Set("Cache-Control", "no-cache")
140160 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)
142162 cmd.Env = append(os.Environ(), gitProtocolEnv(r)...)
143 cmd.Stdin = body
144 cmd.Stdout = w
145 cmd.Run()
163 cmd.Stdin = br
164 cmd.Stdout = out
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))
146240}
147241
148242// 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) {
23402340 s.notFound(w, r)
23412341 return
23422342 }
2343 out, kill, finish, ok := s.packSlot(w, r)
2344 if !ok {
2345 return
2346 }
2347 defer finish()
23432348 prefix := fmt.Sprintf("%s-%s", p.Repo.Name, ref)
23442349 w.Header().Set("Content-Type", "application/gzip")
23452350 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)
23472352}
23482353
23492354func 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
22
33import (
44 "bytes"
5 "io"
6 "os"
57 "path/filepath"
68 "strings"
79 "testing"
10 "time"
811
912 "gitbay.org/gitbay/internal/config"
1013 "gitbay.org/gitbay/internal/control"
14 "gitbay.org/gitbay/internal/gitutil"
15 "gitbay.org/gitbay/internal/packlimit"
1116 "gitbay.org/gitbay/internal/protocol"
1217 "gitbay.org/gitbay/internal/store"
1318)
@@ -51,7 +56,7 @@ func TestRefusedPushIsAudited(t *testing.T) {
5156 bob.Pending = pending
5257 key := store.SSHKey{Scope: "full", Fingerprint: "SHA256:test"}
5358 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",
5560 strings.NewReader(""), &out, &errOut, nil, nil, nil)
5661 if code != protocol.ExitDenied {
5762 t.Fatalf("pending %v: exit %d: %s", pending, code, errOut.String())
@@ -66,3 +71,171 @@ func TestRefusedPushIsAudited(t *testing.T) {
6671 }
6772 }
6873}
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 (
2929 "gitbay.org/gitbay/internal/control"
3030 "gitbay.org/gitbay/internal/gitutil"
3131 "gitbay.org/gitbay/internal/hookd"
32 "gitbay.org/gitbay/internal/packlimit"
3233 "gitbay.org/gitbay/internal/policy"
3334 "gitbay.org/gitbay/internal/protocol"
3435 "gitbay.org/gitbay/internal/store"
@@ -37,6 +38,7 @@ import (
3738type Server struct {
3839 cfg config.Config
3940 st *store.Store
41 packs *packlimit.Limiter
4042 sshCfg *ssh.ServerConfig
4143 authLimiter *rateLimiter
4244 sessions sync.WaitGroup // accepted connections still being served
@@ -68,8 +70,8 @@ func (c *conn) cut() {
6870 c.net.Close()
6971}
7072
71func New(cfg config.Config, st *store.Store) (*Server, error) {
72 s := &Server{cfg: cfg, st: st, authLimiter: newRateLimiter(cfg.Limits.SSHAuthRate, time.Minute), conns: map[*conn]struct{}{}, stopping: make(chan struct{})}
73func New(cfg config.Config, st *store.Store, packs *packlimit.Limiter) (*Server, error) {
74 s := &Server{cfg: cfg, st: st, packs: packs, authLimiter: newRateLimiter(cfg.Limits.SSHAuthRate, time.Minute), conns: map[*conn]struct{}{}, stopping: make(chan struct{})}
7375
7476 sc := &ssh.ServerConfig{
7577 PublicKeyCallback: s.authenticate,
@@ -423,7 +425,7 @@ func (s *Server) runExec(c *conn, sconn *ssh.ServerConn, ch ssh.Channel, term co
423425 return protocol.ExitDenied
424426 }
425427 _ = 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)
427429}
428430
429431// 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 {
459461// Exec runs one SSH exec command line for an authenticated key. It is the
460462// single dispatch path shared by the embedded listener and the system-sshd
461463// 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,
463465 stdin io.Reader, stdout, stderr io.Writer, done, stopping, revoked <-chan struct{}) int {
464466 if user.Disabled {
465467 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,
477479 if user.Pending {
478480 fmt.Fprintln(stderr, "your account is not active yet: verify your email first")
479481 } 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)
481483 }
482484 // A refused push is a refused write, audited like one. runGit
483485 // 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,
515517}
516518
517519// runGit streams a git transport service after access checks.
518func runGit(cfg config.Config, st *store.Store, user store.User, scope string, argv []string,
519 stdin io.Reader, stdout, stderr io.Writer, revoked <-chan struct{}) int {
520func runGit(cfg config.Config, st *store.Store, packs *packlimit.Limiter, user store.User, scope string, argv []string,
521 stdin io.Reader, stdout, stderr io.Writer, done, stopping, revoked <-chan struct{}) int {
520522 service := argv[0]
521523 if len(argv) != 2 {
522524 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
601603 defer st.DeletePushToken(token)
602604 env = append(env, hookd.EnvToken+"="+token)
603605 }
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 {
605672 return protocol.ExitFailure
606673 }
607674 return protocol.ExitOK
internal/sshd/sshd_test.go +2 −2
@@ -64,7 +64,7 @@ func newTestServer(t *testing.T) testServer {
6464
6565 cfg := config.Default()
6666 cfg.Server.Root = root
67 srv, err := New(cfg, st)
67 srv, err := New(cfg, st, nil)
6868 if err != nil {
6969 t.Fatal(err)
7070 }
@@ -245,7 +245,7 @@ func TestUnregisteredKeyMessageNamesFingerprintAndHost(t *testing.T) {
245245 // The settings link keeps the site URL's scheme and port.
246246 cfg.Server.SiteURL = "http://forge.test:8080/"
247247 cfg.Registration.Mode = "open"
248 srv, err := New(cfg, st)
248 srv, err := New(cfg, st, nil)
249249 if err != nil {
250250 t.Fatal(err)
251251 }