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