Commit 973797757f
Verified · cmc
Layout: unified · split
internal/config/config.go +44 −9
| @@ -35,6 +35,16 @@ const ( | |||
| 35 | DefaultPackQueueWait = time.Minute | 35 | DefaultPackQueueWait = time.Minute |
| 36 | ) | 36 | ) |
| 37 | 37 | ||
| 38 | // Push defaults: receive-pack indexes what it is sent, a core each for a | ||
| 39 | // large push, and has its own budget so clones and pushes cannot starve | ||
| 40 | // each other. | ||
| 41 | const ( | ||
| 42 | DefaultPushConcurrency = 2 | ||
| 43 | DefaultPushPerPrincipal = 1 | ||
| 44 | DefaultPushQueue = 16 | ||
| 45 | DefaultPushQueueWait = time.Minute | ||
| 46 | ) | ||
| 47 | |||
| 38 | type Config struct { | 48 | type Config struct { |
| 39 | Server Server `toml:"server"` | 49 | Server Server `toml:"server"` |
| 40 | SSH SSH `toml:"ssh"` | 50 | SSH SSH `toml:"ssh"` |
| @@ -235,12 +245,32 @@ type Limits struct { | |||
| 235 | PackPerPrincipal int `toml:"pack_per_principal"` | 245 | PackPerPrincipal int `toml:"pack_per_principal"` |
| 236 | PackQueue int `toml:"pack_queue"` | 246 | PackQueue int `toml:"pack_queue"` |
| 237 | PackQueueWait string `toml:"pack_queue_wait"` | 247 | PackQueueWait string `toml:"pack_queue_wait"` |
| 248 | // PushConcurrency caps receive-pack running at once over SSH, with | ||
| 249 | // its hooks; PushPerPrincipal caps it per account or deploy key. | ||
| 250 | // PushQueue and PushQueueWait bound the wait for a slot. The counts | ||
| 251 | // read like the pack_* ones: 0 takes the default, negative is off. | ||
| 252 | PushConcurrency int `toml:"push_concurrency"` | ||
| 253 | PushPerPrincipal int `toml:"push_per_principal"` | ||
| 254 | PushQueue int `toml:"push_queue"` | ||
| 255 | PushQueueWait string `toml:"push_queue_wait"` | ||
| 238 | } | 256 | } |
| 239 | 257 | ||
| 240 | // PackLimits resolves the pack_* settings for packlimit.New. A zero | 258 | // PackLimits resolves the pack_* settings for packlimit.New. A zero |
| 241 | // max or per is no bound; an unbounded queue is math.MaxInt, since | 259 | // max or per is no bound; an unbounded queue is math.MaxInt, since |
| 242 | // packlimit reads a zero queue as no queue at all. | 260 | // packlimit reads a zero queue as no queue at all. |
| 243 | func (l Limits) PackLimits() (max, per, queue int, wait time.Duration) { | 261 | func (l Limits) PackLimits() (max, per, queue int, wait time.Duration) { |
| 262 | return resolveLimits(l.PackConcurrency, l.PackPerPrincipal, l.PackQueue, l.PackQueueWait, | ||
| 263 | DefaultPackConcurrency, DefaultPackPerPrincipal, DefaultPackQueue, DefaultPackQueueWait) | ||
| 264 | } | ||
| 265 | |||
| 266 | // PushLimits resolves the push_* settings the way PackLimits does the | ||
| 267 | // pack_* ones. | ||
| 268 | func (l Limits) PushLimits() (max, per, queue int, wait time.Duration) { | ||
| 269 | return resolveLimits(l.PushConcurrency, l.PushPerPrincipal, l.PushQueue, l.PushQueueWait, | ||
| 270 | DefaultPushConcurrency, DefaultPushPerPrincipal, DefaultPushQueue, DefaultPushQueueWait) | ||
| 271 | } | ||
| 272 | |||
| 273 | func resolveLimits(maxV, perV, queueV int, waitV string, maxDef, perDef, queueDef int, waitDef time.Duration) (max, per, queue int, wait time.Duration) { | ||
| 244 | pick := func(v, def int) int { | 274 | pick := func(v, def int) int { |
| 245 | switch { | 275 | switch { |
| 246 | case v == 0: | 276 | case v == 0: |
| @@ -250,16 +280,15 @@ func (l Limits) PackLimits() (max, per, queue int, wait time.Duration) { | |||
| 250 | } | 280 | } |
| 251 | return v | 281 | return v |
| 252 | } | 282 | } |
| 253 | queue = pick(l.PackQueue, DefaultPackQueue) | 283 | queue = pick(queueV, queueDef) |
| 254 | if l.PackQueue < 0 { | 284 | if queueV < 0 { |
| 255 | queue = math.MaxInt | 285 | queue = math.MaxInt |
| 256 | } | 286 | } |
| 257 | wait = DefaultPackQueueWait | 287 | wait = waitDef |
| 258 | if d, err := time.ParseDuration(l.PackQueueWait); err == nil && d > 0 { | 288 | if d, err := time.ParseDuration(waitV); err == nil && d > 0 { |
| 259 | wait = d | 289 | wait = d |
| 260 | } | 290 | } |
| 261 | return pick(l.PackConcurrency, DefaultPackConcurrency), | 291 | return pick(maxV, maxDef), pick(perV, perDef), queue, wait |
| 262 | pick(l.PackPerPrincipal, DefaultPackPerPrincipal), queue, wait | ||
| 263 | } | 292 | } |
| 264 | 293 | ||
| 265 | type Mail struct { | 294 | type Mail struct { |
| @@ -619,9 +648,15 @@ func (c Config) Validate() error { | |||
| 619 | if c.Limits.MaxReposPerUser < 0 || c.Limits.MaxBytesPerUser < 0 || c.Limits.MaxSnippetsPerUser < 0 { | 648 | if c.Limits.MaxReposPerUser < 0 || c.Limits.MaxBytesPerUser < 0 || c.Limits.MaxSnippetsPerUser < 0 { |
| 620 | errs = append(errs, errors.New("limits.max_repos_per_user, max_bytes_per_user and max_snippets_per_user must not be negative")) | 649 | errs = append(errs, errors.New("limits.max_repos_per_user, max_bytes_per_user and max_snippets_per_user must not be negative")) |
| 621 | } | 650 | } |
| 622 | if w := c.Limits.PackQueueWait; w != "" { | 651 | for _, w := range []struct{ name, val string }{ |
| 623 | if d, err := time.ParseDuration(w); err != nil || d <= 0 { | 652 | {"pack_queue_wait", c.Limits.PackQueueWait}, |
| 624 | errs = append(errs, fmt.Errorf("limits.pack_queue_wait %q must be a positive duration such as 60s", w)) | 653 | {"push_queue_wait", c.Limits.PushQueueWait}, |
| 654 | } { | ||
| 655 | if w.val == "" { | ||
| 656 | continue | ||
| 657 | } | ||
| 658 | if d, err := time.ParseDuration(w.val); err != nil || d <= 0 { | ||
| 659 | errs = append(errs, fmt.Errorf("limits.%s %q must be a positive duration such as 60s", w.name, w.val)) | ||
| 625 | } | 660 | } |
| 626 | } | 661 | } |
| 627 | if c.Push.Enabled { | 662 | if c.Push.Enabled { |
internal/config/config_test.go +32
| @@ -63,6 +63,33 @@ func TestPackLimits(t *testing.T) { | |||
| 63 | } | 63 | } |
| 64 | } | 64 | } |
| 65 | 65 | ||
| 66 | // push_* resolve like pack_*, from their own defaults. | ||
| 67 | func TestPushLimits(t *testing.T) { | ||
| 68 | max, per, queue, wait := Limits{}.PushLimits() | ||
| 69 | if max != DefaultPushConcurrency || per != DefaultPushPerPrincipal || queue != DefaultPushQueue || wait != DefaultPushQueueWait { | ||
| 70 | t.Fatalf("defaults: %d %d %d %s", max, per, queue, wait) | ||
| 71 | } | ||
| 72 | if max != 2 || per != 1 || queue != 16 || wait != time.Minute { | ||
| 73 | t.Fatalf("defaults moved: %d %d %d %s", max, per, queue, wait) | ||
| 74 | } | ||
| 75 | max, per, queue, wait = Limits{PushConcurrency: -1, PushPerPrincipal: -1, PushQueue: -1, PushQueueWait: "5s"}.PushLimits() | ||
| 76 | if max != 0 || per != 0 || queue != math.MaxInt || wait != 5*time.Second { | ||
| 77 | t.Fatalf("off: %d %d %d %s", max, per, queue, wait) | ||
| 78 | } | ||
| 79 | cfg, err := Load(writeConfig(t, minimal+"\n[limits]\npush_concurrency = 4\npush_per_principal = 2\npush_queue = 8\npush_queue_wait = \"30s\"\n")) | ||
| 80 | if err != nil { | ||
| 81 | t.Fatal(err) | ||
| 82 | } | ||
| 83 | max, per, queue, wait = cfg.Limits.PushLimits() | ||
| 84 | if max != 4 || per != 2 || queue != 8 || wait != 30*time.Second { | ||
| 85 | t.Fatalf("loaded: %d %d %d %s", max, per, queue, wait) | ||
| 86 | } | ||
| 87 | // The pack budget is not read from the push settings. | ||
| 88 | if pm, _, _, _ := cfg.Limits.PackLimits(); pm != DefaultPackConcurrency { | ||
| 89 | t.Fatalf("pack_concurrency %d, want the default", pm) | ||
| 90 | } | ||
| 91 | } | ||
| 92 | |||
| 66 | func TestContradictions(t *testing.T) { | 93 | func TestContradictions(t *testing.T) { |
| 67 | cases := []struct { | 94 | cases := []struct { |
| 68 | name string | 95 | name string |
| @@ -74,6 +101,11 @@ func TestContradictions(t *testing.T) { | |||
| 74 | minimal + "\n[limits]\npack_queue_wait = \"soon\"\n", | 101 | minimal + "\n[limits]\npack_queue_wait = \"soon\"\n", |
| 75 | "limits.pack_queue_wait", | 102 | "limits.pack_queue_wait", |
| 76 | }, | 103 | }, |
| 104 | { | ||
| 105 | "bad push_queue_wait", | ||
| 106 | minimal + "\n[limits]\npush_queue_wait = \"-5s\"\n", | ||
| 107 | "limits.push_queue_wait", | ||
| 108 | }, | ||
| 77 | { | 109 | { |
| 78 | "registration open without smtp", | 110 | "registration open without smtp", |
| 79 | minimal + "\n[registration]\nmode = \"open\"\n", | 111 | minimal + "\n[registration]\nmode = \"open\"\n", |
internal/packlimit/packlimit.go +18 −7
| @@ -1,8 +1,9 @@ | |||
| 1 | // Package packlimit bounds concurrent git pack generation. upload-pack | 1 | // Package packlimit bounds concurrent git processes. upload-pack and |
| 2 | // and upload-archive over SSH, smart HTTP and git:// draw on one | 2 | // upload-archive over SSH, smart HTTP and git:// draw on one budget, |
| 3 | // budget: a global cap, a cap per principal (an account, or a client | 3 | // receive-pack on a second: each a global cap, a cap per principal (an |
| 4 | // address on the anonymous transports), and a bounded queue whose | 4 | // account, a deploy key, or a client address on the anonymous |
| 5 | // waiters give up after a fixed wait or when the client goes away. | 5 | // transports), and a bounded queue whose waiters give up after a fixed |
| 6 | // wait or when the client goes away. | ||
| 6 | // Waiters are not served in order; a new arrival can take a freed slot | 7 | // Waiters are not served in order; a new arrival can take a freed slot |
| 7 | // ahead of them, and the wait bounds how long any one of them waits. | 8 | // ahead of them, and the wait bounds how long any one of them waits. |
| 8 | package packlimit | 9 | package packlimit |
| @@ -24,6 +25,7 @@ var ( | |||
| 24 | type Limiter struct { | 25 | type Limiter struct { |
| 25 | max, per, queue int | 26 | max, per, queue int |
| 26 | wait time.Duration | 27 | wait time.Duration |
| 28 | name string // what is limited, for the refusal log | ||
| 27 | 29 | ||
| 28 | // Principals starting with class may hold at most classCap slots | 30 | // Principals starting with class may hold at most classCap slots |
| 29 | // between them; classCap 0 is no class cap. | 31 | // between them; classCap 0 is no class cap. |
| @@ -45,7 +47,7 @@ func New(max, per, queue int, wait time.Duration) *Limiter { | |||
| 45 | if max <= 0 { | 47 | if max <= 0 { |
| 46 | return nil | 48 | return nil |
| 47 | } | 49 | } |
| 48 | return &Limiter{max: max, per: per, queue: queue, wait: wait, | 50 | return &Limiter{max: max, per: per, queue: queue, wait: wait, name: "pack", |
| 49 | held: map[string]int{}, waiting: map[string]int{}, changed: make(chan struct{}), | 51 | held: map[string]int{}, waiting: map[string]int{}, changed: make(chan struct{}), |
| 50 | warned: map[string]time.Time{}} | 52 | warned: map[string]time.Time{}} |
| 51 | } | 53 | } |
| @@ -71,10 +73,19 @@ func (l *Limiter) Refused(transport, principal string, err error) { | |||
| 71 | if errors.Is(err, ErrGone) { | 73 | if errors.Is(err, ErrGone) { |
| 72 | reason = "gone" | 74 | reason = "gone" |
| 73 | } | 75 | } |
| 74 | slog.Warn("pack limit: request turned away (logged at most once a minute per transport)", | 76 | slog.Warn(l.name+" limit: request turned away (logged at most once a minute per transport)", |
| 75 | "transport", transport, "class", class, "reason", reason) | 77 | "transport", transport, "class", class, "reason", reason) |
| 76 | } | 78 | } |
| 77 | 79 | ||
| 80 | // Name sets what the refusal log calls this limit ("pack" unless set). | ||
| 81 | // Call it before the limiter is in use. | ||
| 82 | func (l *Limiter) Name(name string) { | ||
| 83 | if l == nil { | ||
| 84 | return | ||
| 85 | } | ||
| 86 | l.name = name | ||
| 87 | } | ||
| 88 | |||
| 78 | // CapClass caps the slots that principals starting with prefix may hold | 89 | // CapClass caps the slots that principals starting with prefix may hold |
| 79 | // between them. Call it before the limiter is in use. | 90 | // between them. Call it before the limiter is in use. |
| 80 | func (l *Limiter) CapClass(prefix string, n int) { | 91 | func (l *Limiter) CapClass(prefix string, n int) { |
internal/packlimit/packlimit_test.go +24
| @@ -313,3 +313,27 @@ func TestRefusedLogsOncePerTransport(t *testing.T) { | |||
| 313 | var none *Limiter | 313 | var none *Limiter |
| 314 | none.Refused("git", "ip:x", ErrBusy) | 314 | none.Refused("git", "ip:x", ErrBusy) |
| 315 | } | 315 | } |
| 316 | |||
| 317 | // Two limiters log apart: a named one says what it limits, and its | ||
| 318 | // once-a-minute window does not silence the other's. | ||
| 319 | func TestRefusedNamesTheLimit(t *testing.T) { | ||
| 320 | var buf bytes.Buffer | ||
| 321 | old := slog.Default() | ||
| 322 | slog.SetDefault(slog.New(slog.NewTextHandler(&buf, nil))) | ||
| 323 | t.Cleanup(func() { slog.SetDefault(old) }) | ||
| 324 | |||
| 325 | packs := New(1, 0, 0, time.Second) | ||
| 326 | pushes := New(1, 0, 0, time.Second) | ||
| 327 | pushes.Name("push") | ||
| 328 | packs.Refused("ssh", "user:4", ErrBusy) | ||
| 329 | pushes.Refused("ssh", "key:9", ErrBusy) | ||
| 330 | out := buf.String() | ||
| 331 | if !strings.Contains(out, `"pack limit: request turned away`) || !strings.Contains(out, `"push limit: request turned away`) { | ||
| 332 | t.Fatalf("want one line per limit:\n%s", out) | ||
| 333 | } | ||
| 334 | if !strings.Contains(out, "class=key") || strings.Contains(out, "key:9") { | ||
| 335 | t.Fatalf("deploy key principal logged or class missing:\n%s", out) | ||
| 336 | } | ||
| 337 | var none *Limiter | ||
| 338 | none.Name("push") | ||
| 339 | } | ||