internal/sshd/sshd.go
802 lines · 25331 bytes
1// Package sshd implements the embedded SSH listener: public-key auth against
2// registered keys, then dispatch to git transport or control commands.
3package sshd
4
5import (
6 "bytes"
7 "context"
8 "crypto/ed25519"
9 "crypto/rand"
10 "encoding/base64"
11 "encoding/pem"
12 "errors"
13 "fmt"
14 "io"
15 "log/slog"
16 "maps"
17 "net"
18 "os"
19 "path/filepath"
20 "slices"
21 "strconv"
22 "strings"
23 "sync"
24 "sync/atomic"
25 "time"
26
27 "golang.org/x/crypto/ssh"
28
29 "gitbay.org/gitbay/internal/config"
30 "gitbay.org/gitbay/internal/control"
31 "gitbay.org/gitbay/internal/gitutil"
32 "gitbay.org/gitbay/internal/hookd"
33 "gitbay.org/gitbay/internal/packlimit"
34 "gitbay.org/gitbay/internal/policy"
35 "gitbay.org/gitbay/internal/protocol"
36 "gitbay.org/gitbay/internal/store"
37)
38
39type Server struct {
40 cfg config.Config
41 st *store.Store
42 packs *packlimit.Limiter
43 pushes *packlimit.Limiter
44 sshCfg *ssh.ServerConfig
45 authLimiter *rateLimiter
46 sessions sync.WaitGroup // accepted connections still being served
47 mu sync.Mutex
48 conns map[*conn]struct{}
49 stopping chan struct{} // closed by Stop
50 stopOnce sync.Once
51}
52
53// conn is one accepted connection and how many sessions it is running.
54// A CLI's shared connection sits idle between commands; on shutdown an
55// idle connection is closed at once and only a session mid-command is
56// waited for (#141).
57type conn struct {
58 net net.Conn
59 active atomic.Int32
60 // keyID and userID are the key that authenticated the connection and
61 // its account: 0 before the handshake and for an unregistered key.
62 // Guarded by Server.mu.
63 keyID, userID int64
64 revoked chan struct{} // closed by cut
65 cutOnce sync.Once
66}
67
68// cut ends the connection because its key was revoked: a git transport
69// on it is killed, and every other command loses its channel.
70func (c *conn) cut() {
71 c.cutOnce.Do(func() { close(c.revoked) })
72 c.net.Close()
73}
74
75func New(cfg config.Config, st *store.Store, packs, pushes *packlimit.Limiter) (*Server, error) {
76 s := &Server{cfg: cfg, st: st, packs: packs, pushes: pushes, authLimiter: newRateLimiter(cfg.Limits.SSHAuthRate, time.Minute), conns: map[*conn]struct{}{}, stopping: make(chan struct{})}
77
78 sc := &ssh.ServerConfig{
79 PublicKeyCallback: s.authenticate,
80 ServerVersion: "SSH-2.0-gitbayd",
81 }
82 signers, err := loadHostKeys(cfg)
83 if err != nil {
84 return nil, err
85 }
86 for _, sg := range signers {
87 sc.AddHostKey(sg)
88 }
89 s.sshCfg = sc
90 st.OnRevoke(s.revoke)
91 return s, nil
92}
93
94// loadHostKeys loads the configured host keys, or generates an ed25519 key
95// under server.root/ssh/ when none are configured.
96func loadHostKeys(cfg config.Config) ([]ssh.Signer, error) {
97 paths := cfg.SSH.HostKeys
98 if len(paths) == 0 {
99 p := filepath.Join(cfg.Server.Root, "ssh", "host_ed25519")
100 if _, err := os.Stat(p); errors.Is(err, os.ErrNotExist) {
101 if err := generateHostKey(p); err != nil {
102 return nil, fmt.Errorf("generating host key: %w", err)
103 }
104 slog.Info("generated ssh host key", "path", p)
105 }
106 paths = []string{p}
107 }
108 var signers []ssh.Signer
109 for _, p := range paths {
110 raw, err := os.ReadFile(p)
111 if err != nil {
112 return nil, fmt.Errorf("host key %s: %w", p, err)
113 }
114 sg, err := ssh.ParsePrivateKey(raw)
115 if err != nil {
116 return nil, fmt.Errorf("host key %s: %w", p, err)
117 }
118 signers = append(signers, sg)
119 }
120 return signers, nil
121}
122
123func generateHostKey(path string) error {
124 if err := os.MkdirAll(filepath.Dir(path), 0o700); err != nil {
125 return err
126 }
127 _, priv, err := ed25519.GenerateKey(rand.Reader)
128 if err != nil {
129 return err
130 }
131 block, err := ssh.MarshalPrivateKey(priv, "")
132 if err != nil {
133 return err
134 }
135 return os.WriteFile(path, pem.EncodeToMemory(block), 0o600)
136}
137
138// authenticate resolves the presented key to a registered account. The SSH
139// username is ignored; identity comes from the key alone. When registration
140// is open or invite-based, unknown keys are admitted to run exactly one
141// command: register.
142func (s *Server) authenticate(meta ssh.ConnMetadata, pub ssh.PublicKey) (*ssh.Permissions, error) {
143 ip := remoteIP(meta.RemoteAddr())
144 if !s.authLimiter.allow(ip) {
145 // One audit entry per throttled window, not per rejected attempt.
146 if s.authLimiter.firstThrottle(ip) {
147 s.st.Audit(0, "auth.throttled", map[string]any{"ip": ip, "rate": s.cfg.Limits.SSHAuthRate})
148 }
149 return nil, fmt.Errorf("too many authentication attempts; try again shortly")
150 }
151 fp := ssh.FingerprintSHA256(pub)
152 key, err := s.st.SSHKeyByFingerprint(fp)
153 if err != nil && !errors.Is(err, store.ErrNotFound) {
154 // The store, not the key, failed. Neither a failure against the
155 // limiter nor "unknown key": a busy database during a restart
156 // would otherwise lock every client out for a minute.
157 slog.Error("ssh auth: key lookup", "err", err)
158 return nil, fmt.Errorf("authentication temporarily unavailable")
159 }
160 if err != nil {
161 if s.cfg.Registration.Mode != "closed" {
162 return &ssh.Permissions{Extensions: map[string]string{
163 "anon-key": base64.StdEncoding.EncodeToString(pub.Marshal()),
164 }}, nil
165 }
166 s.authLimiter.fail(ip)
167 s.st.Audit(0, "auth.failed", map[string]any{"ip": ip, "fingerprint": fp})
168 return nil, fmt.Errorf("unknown key %s", fp)
169 }
170 if key.Expired(time.Now()) {
171 s.authLimiter.fail(ip)
172 s.st.Audit(key.UserID, "auth.expired", map[string]any{"ip": ip, "fingerprint": fp})
173 return nil, fmt.Errorf("key %s has expired", fp)
174 }
175 s.authLimiter.success(ip)
176 return &ssh.Permissions{Extensions: map[string]string{
177 "user-id": strconv.FormatInt(key.UserID, 10),
178 "key-id": strconv.FormatInt(key.ID, 10),
179 }}, nil
180}
181
182// Serve accepts connections on ln until it is closed.
183func (s *Server) Serve(ln net.Listener) error {
184 served := make(chan struct{})
185 defer close(served)
186 go s.sweep(served)
187 for {
188 nc, err := ln.Accept()
189 if err != nil {
190 return err
191 }
192 c := &conn{net: nc, revoked: make(chan struct{})}
193 s.mu.Lock()
194 s.conns[c] = struct{}{}
195 s.mu.Unlock()
196 s.sessions.Add(1)
197 go func() {
198 defer s.sessions.Done()
199 defer func() {
200 s.mu.Lock()
201 delete(s.conns, c)
202 s.mu.Unlock()
203 }()
204 s.handleConn(c)
205 }()
206 }
207}
208
209// revoke closes the connections opened by the keys r names.
210func (s *Server) revoke(r store.Revoked) {
211 var cut []*conn
212 s.mu.Lock()
213 for c := range s.conns {
214 if c.keyID == 0 {
215 continue
216 }
217 if (r.UserID != 0 && c.userID == r.UserID) || slices.Contains(r.KeyIDs, c.keyID) {
218 cut = append(cut, c)
219 }
220 }
221 s.mu.Unlock()
222 for _, c := range cut {
223 c.cut()
224 }
225}
226
227// sweepInterval bounds how long a revocation this process was not told
228// about (gitbayd admin on the host) leaves a connection open.
229const sweepInterval = 15 * time.Second
230
231func (s *Server) sweep(served <-chan struct{}) {
232 t := time.NewTicker(sweepInterval)
233 defer t.Stop()
234 for {
235 select {
236 case <-t.C:
237 s.sweepOnce()
238 case <-served:
239 return
240 case <-s.stopping:
241 return
242 }
243 }
244}
245
246// sweepOnce cuts every connection whose key is no longer live. Only
247// connections whose key was asked about are judged: one that
248// authenticated while the query ran waits for the next sweep.
249func (s *Server) sweepOnce() {
250 asked := map[int64]bool{}
251 s.mu.Lock()
252 for c := range s.conns {
253 if c.keyID != 0 {
254 asked[c.keyID] = true
255 }
256 }
257 s.mu.Unlock()
258 if len(asked) == 0 {
259 return
260 }
261 live, err := s.st.LiveSSHKeys(slices.Collect(maps.Keys(asked)))
262 if err != nil {
263 slog.Error("ssh sweep: key lookup", "err", err)
264 return
265 }
266 var cut []*conn
267 s.mu.Lock()
268 for c := range s.conns {
269 if asked[c.keyID] && !live[c.keyID] {
270 cut = append(cut, c)
271 }
272 }
273 s.mu.Unlock()
274 for _, c := range cut {
275 c.cut()
276 }
277}
278
279// Stop ends the commands that run until something happens (build log
280// --follow), so a shutdown drain waits only for work that finishes. It
281// does not close connections; Shutdown does.
282func (s *Server) Stop() {
283 s.stopOnce.Do(func() { close(s.stopping) })
284}
285
286// Shutdown closes every idle connection, then waits for the ones with a
287// session running, or for ctx. The caller closes the listener first; a
288// push in flight completes rather than being cut mid-pack.
289func (s *Server) Shutdown(ctx context.Context) error {
290 s.Stop()
291 s.mu.Lock()
292 for c := range s.conns {
293 if c.active.Load() == 0 {
294 c.net.Close()
295 }
296 }
297 s.mu.Unlock()
298 done := make(chan struct{})
299 go func() {
300 s.sessions.Wait()
301 close(done)
302 }()
303 select {
304 case <-done:
305 return nil
306 case <-ctx.Done():
307 return ctx.Err()
308 }
309}
310
311func (s *Server) handleConn(c *conn) {
312 defer c.net.Close()
313 sconn, chans, reqs, err := ssh.NewServerConn(c.net, s.sshCfg)
314 if err != nil {
315 return
316 }
317 defer sconn.Close()
318 ext := sconn.Permissions.Extensions
319 s.mu.Lock()
320 c.keyID, _ = strconv.ParseInt(ext["key-id"], 10, 64)
321 c.userID, _ = strconv.ParseInt(ext["user-id"], 10, 64)
322 s.mu.Unlock()
323 go ssh.DiscardRequests(reqs)
324
325 for newCh := range chans {
326 if newCh.ChannelType() != "session" {
327 newCh.Reject(ssh.UnknownChannelType, "only session channels are supported")
328 continue
329 }
330 ch, chReqs, err := newCh.Accept()
331 if err != nil {
332 continue
333 }
334 c.active.Add(1)
335 go func() {
336 defer c.active.Add(-1)
337 s.handleSession(c, sconn, ch, chReqs)
338 }()
339 }
340}
341
342func (s *Server) handleSession(c *conn, sconn *ssh.ServerConn, ch ssh.Channel, reqs <-chan *ssh.Request) {
343 defer ch.Close()
344 var term control.Term
345 for req := range reqs {
346 switch req.Type {
347 case "exec":
348 var payload struct{ Command string }
349 if err := ssh.Unmarshal(req.Payload, &payload); err != nil {
350 req.Reply(false, nil)
351 continue
352 }
353 req.Reply(true, nil)
354 // x/crypto closes reqs when the client closes the channel. That
355 // is how a follow learns nobody is reading: the CLI's shared
356 // connection outlives a Ctrl-C, the channel does not. Stop
357 // ends it too, for a restart.
358 closed := make(chan struct{})
359 go func() {
360 for r := range reqs {
361 r.Reply(false, nil)
362 }
363 close(closed)
364 }()
365 done := make(chan struct{})
366 go func() {
367 select {
368 case <-closed:
369 case <-s.stopping:
370 }
371 close(done)
372 }()
373 code := s.runExec(c, sconn, ch, term, payload.Command, done)
374 sendExit(ch, code)
375 return
376 case "shell":
377 req.Reply(true, nil)
378 fmt.Fprintf(ch, "gitbay control plane: interactive shells are not available.\nTry: ssh %s help\n", s.cfg.Server.SiteURL)
379 sendExit(ch, protocol.ExitUsage)
380 return
381 case "env":
382 var kv struct{ Name, Value string }
383 if ssh.Unmarshal(req.Payload, &kv) == nil && kv.Name == "GITBAY_TERM" {
384 term = control.ParseTerm(kv.Value)
385 }
386 req.Reply(true, nil)
387 case "pty-req":
388 // Harmless; accept and ignore.
389 req.Reply(true, nil)
390 default:
391 req.Reply(false, nil)
392 }
393 }
394}
395
396func sendExit(ch ssh.Channel, code int) {
397 var msg = struct{ Status uint32 }{uint32(code)}
398 ch.SendRequest("exit-status", false, ssh.Marshal(&msg))
399}
400
401func (s *Server) runExec(c *conn, sconn *ssh.ServerConn, ch ssh.Channel, term control.Term, cmdline string, done <-chan struct{}) int {
402 ext := sconn.Permissions.Extensions
403 if blob := ext["anon-key"]; blob != "" {
404 return s.runAnonymous(ch, blob, cmdline)
405 }
406 userID, _ := strconv.ParseInt(ext["user-id"], 10, 64)
407 keyID, _ := strconv.ParseInt(ext["key-id"], 10, 64)
408 // A connection outlives its commands, so the key is read again for
409 // each one: what it may do is what it may do now (#256).
410 key, err := s.st.SSHKeyByID(keyID)
411 if errors.Is(err, store.ErrNotFound) || (err == nil && key.UserID != userID) {
412 fmt.Fprintln(ch.Stderr(), "this key is no longer registered")
413 return protocol.ExitDenied
414 }
415 if err != nil {
416 slog.Error("ssh exec: key lookup", "err", err)
417 fmt.Fprintln(ch.Stderr(), "authentication temporarily unavailable")
418 return protocol.ExitFailure
419 }
420 if key.Expired(time.Now()) {
421 fmt.Fprintln(ch.Stderr(), "this key has expired; remove it and add a new one")
422 return protocol.ExitDenied
423 }
424 user, err := s.st.UserByID(userID)
425 if err != nil {
426 fmt.Fprintln(ch.Stderr(), "account no longer exists")
427 return protocol.ExitDenied
428 }
429 _ = s.st.TouchSSHKey(keyID)
430 return Exec(s.cfg, s.st, s.packs, s.pushes, user, key, term, cmdline, ch, ch, ch.Stderr(), done, s.stopping, c.revoked)
431}
432
433// runAnonymous handles a session from an unregistered key: the register
434// command and nothing else.
435func (s *Server) runAnonymous(ch ssh.Channel, keyB64, cmdline string) int {
436 raw, err := base64.StdEncoding.DecodeString(keyB64)
437 if err != nil {
438 return protocol.ExitFailure
439 }
440 pub, err := ssh.ParsePublicKey(raw)
441 if err != nil {
442 return protocol.ExitFailure
443 }
444 argv, err := protocol.Tokenize(cmdline)
445 if err != nil {
446 fmt.Fprintf(ch.Stderr(), "cannot parse command: %v\n", err)
447 return protocol.ExitUsage
448 }
449 if len(argv) == 0 || argv[0] != "register" {
450 host := s.cfg.SiteHost()
451 fp := ssh.FingerprintSHA256(pub)
452 flag := map[string]string{"open": "--email <address>", "invite": "--invite <code>"}[s.cfg.Registration.Mode]
453 fmt.Fprintf(ch.Stderr(),
454 "this key (%s) is not registered on %s.\n"+
455 "already have an account? add it at %s/settings#keys\n"+
456 "new here? ssh git@%s register --username <name> %s\n",
457 fp, host, strings.TrimSuffix(s.cfg.Server.SiteURL, "/"), host, flag)
458 return protocol.ExitDenied
459 }
460 return control.RunRegister(s.cfg, s.st, pub, argv, ch, ch.Stderr())
461}
462
463// Exec runs one SSH exec command line for an authenticated key. It is the
464// single dispatch path shared by the embedded listener and the system-sshd
465// forced command (gitbayd shell). Closing revoked kills a git transport.
466// packs bounds clones and fetches, and repo download; pushes bounds
467// receive-pack. A nil limiter is no limit.
468func Exec(cfg config.Config, st *store.Store, packs, pushes *packlimit.Limiter, user store.User, key store.SSHKey, term control.Term, cmdline string,
469 stdin io.Reader, stdout, stderr io.Writer, done, stopping, revoked <-chan struct{}) int {
470 // A person signing in with a full-scope key cancels a scheduled
471 // deletion; automation on narrower keys is refused and cannot.
472 if user.Disabled && user.DeleteAfter != "" && key.Scope == "full" && control.CancelScheduledDeletion(st, &user, "ssh") {
473 fmt.Fprintln(stderr, "deletion of your account was cancelled")
474 }
475 if user.Disabled {
476 if user.DeleteAfter != "" {
477 fmt.Fprintln(stderr, control.ScheduledRefusal(user))
478 return protocol.ExitDenied
479 }
480 fmt.Fprintln(stderr, "this account is disabled; contact the instance admin")
481 return protocol.ExitDenied
482 }
483 argv, err := protocol.Tokenize(cmdline)
484 if err != nil {
485 fmt.Fprintf(stderr, "cannot parse command: %v\n", err)
486 return protocol.ExitUsage
487 }
488 if len(argv) > 0 {
489 switch argv[0] {
490 case "git-upload-pack", "git-receive-pack", "git-upload-archive":
491 code := protocol.ExitDenied
492 if user.Pending {
493 fmt.Fprintln(stderr, "your account is not active yet: verify your email first")
494 } else {
495 code = runGit(cfg, st, packs, pushes, user, key, argv, stdin, stdout, stderr, done, stopping, revoked)
496 }
497 // A refused push is a refused write, audited like one. runGit
498 // refuses only with the path as the one argument, so argv[1:]
499 // holds no value beyond the target.
500 if argv[0] == "git-receive-pack" && (code == protocol.ExitDenied || code == protocol.ExitNotFound) {
501 control.AuditRefused(st, user.ID, "refused git-receive-pack",
502 map[string]any{"argv": argv[1:], "source": key.Fingerprint, "exit": code})
503 }
504 return code
505 case "git-lfs-authenticate":
506 // Part of the git transport, not the control plane: usable by
507 // git-scoped and deploy keys, with the transports' access rules.
508 if user.Pending {
509 fmt.Fprintln(stderr, "your account is not active yet: verify your email first")
510 return protocol.ExitDenied
511 }
512 return runLFSAuthenticate(cfg, st, user, key, argv, stdout, stderr)
513 }
514 }
515 ctx := &control.Ctx{
516 User: user,
517 Scope: key.Scope,
518 Source: key.Fingerprint,
519 Term: term,
520 Store: st,
521 Cfg: cfg,
522 Stdin: stdin,
523 Stdout: stdout,
524 Stderr: stderr,
525 Done: done,
526 Stopping: stopping,
527 Expires: key.ExpiresAt,
528 Packs: packs,
529 }
530 return control.Dispatch(ctx, argv)
531}
532
533// runGit streams a git transport service after access checks.
534func runGit(cfg config.Config, st *store.Store, packs, pushes *packlimit.Limiter, user store.User, key store.SSHKey, argv []string,
535 stdin io.Reader, stdout, stderr io.Writer, done, stopping, revoked <-chan struct{}) int {
536 service, scope := argv[0], key.Scope
537 if len(argv) != 2 {
538 fmt.Fprintf(stderr, "usage: %s <path>\n", service)
539 return protocol.ExitUsage
540 }
541 write := service == "git-receive-pack"
542 cancel, keepAlive := revoked, 0
543
544 repo, err := st.RepoByPath(argv[1])
545 if err != nil {
546 fmt.Fprintln(stderr, "repository not found")
547 return protocol.ExitNotFound
548 }
549 if policy.IsDeployScope(scope) {
550 // A deploy key authorizes by its binding alone: one repository,
551 // its mode, nothing inherited from whoever registered it. Any
552 // mismatch reads as nonexistence, same as the access rules.
553 if !policy.DeployScopeAllows(scope, repo.ID, write) {
554 fmt.Fprintln(stderr, "repository not found")
555 return protocol.ExitNotFound
556 }
557 } else {
558 grant, err := st.AccessRole(repo.ID, user.ID)
559 if err != nil {
560 fmt.Fprintln(stderr, "internal error")
561 return protocol.ExitFailure
562 }
563 if !policy.CanRead(user, repo, grant) {
564 // Same answer as nonexistence: private repos must not be enumerable.
565 fmt.Fprintln(stderr, "repository not found")
566 return protocol.ExitNotFound
567 }
568 if !policy.ScopeAllowsGit(scope, repo.Path(), write) {
569 fmt.Fprintf(stderr, "this key's scope (%s) does not allow %s on %s\n", scope, service, repo.Path())
570 return protocol.ExitDenied
571 }
572 if write && !policy.CanWrite(user, repo, grant) {
573 fmt.Fprintf(stderr, "write access to %s denied\n", repo.Path())
574 return protocol.ExitDenied
575 }
576 }
577 if write && repo.Settings.Archived {
578 fmt.Fprintf(stderr, "%s is archived and read-only\n", repo.Path())
579 return protocol.ExitDenied
580 }
581 if write {
582 if mirrored, err := st.PullMirrored(repo.ID); err == nil && mirrored {
583 fmt.Fprintf(stderr, "%s is a pull mirror: its refs come from the upstream; push there instead\n", repo.Path())
584 return protocol.ExitDenied
585 }
586 }
587
588 dir := control.RepoDir(cfg.Server.Root, repo.OwnerName, repo.Name)
589 env := []string{
590 hookd.EnvSocket + "=" + hookd.SocketPath(cfg.Server.Root),
591 hookd.EnvRepoID + "=" + strconv.FormatInt(repo.ID, 10),
592 hookd.EnvUserID + "=" + strconv.FormatInt(user.ID, 10),
593 hookd.EnvScope + "=" + scope,
594 }
595 // A storage quota on the owner rides the same mechanism as the pack
596 // cap: the pack may be no larger than what the owner has left.
597 maxPack := cfg.Limits.MaxPackBytes
598 if write {
599 if limit := control.ByteLimit(st, control.QuotaConfig(cfg), repo.OwnerKind, repo.OwnerID); limit > 0 {
600 used := control.OwnedBytes(st, cfg.Server.Root, repo.OwnerKind, repo.OwnerID)
601 left := limit - used
602 if left <= 0 {
603 fmt.Fprintf(stderr, "%s's storage quota is used up (%d of %d bytes); delete something, or ask an admin to raise the limit\n", repo.OwnerName, used, limit)
604 return protocol.ExitDenied
605 }
606 if maxPack == 0 || left < maxPack {
607 maxPack = left
608 }
609 }
610 }
611 if write {
612 // Pushes have their own budget, so a clone storm cannot starve
613 // them or the reverse. The slot covers receive-pack and both
614 // hooks: git waits for post-receive (RefsUpdated) before it
615 // exits, and that work is the push's cost. A deploy key is its
616 // own principal, not the account that registered it.
617 principal := "user:" + strconv.FormatInt(user.ID, 10)
618 if policy.IsDeployScope(scope) {
619 principal = "key:" + strconv.FormatInt(key.ID, 10)
620 }
621 release, code := takeSlot(pushes, principal, done, stderr,
622 "the server is busy: it is at its limit of concurrent pushes; try again in a minute")
623 if code != protocol.ExitOK {
624 return code
625 }
626 slotAt := time.Now()
627 // Deferred before Transport runs, so it fires after receive-pack
628 // has exited, on every path: the client hanging up ends its
629 // stdin and receive-pack with it, and a revoked key kills it.
630 defer release()
631 // hookd answers only a hook that names this receive-pack.
632 token, err := st.CreatePushToken(repo.ID, user.ID, scope)
633 if err != nil {
634 fmt.Fprintln(stderr, "internal error")
635 return protocol.ExitFailure
636 }
637 defer st.DeletePushToken(token)
638 env = append(env, hookd.EnvToken+"="+token)
639 idleFor, receiveFor := cfg.Limits.PushTimeouts()
640 keepAlive = pushKeepAlive(idleFor)
641 if pushes != nil {
642 // A client that holds its slot while sending nothing, or
643 // trickles its pack, is cut: when pre-receive has not
644 // started receiveFor after the slot was taken, or after
645 // idleFor with no byte either way once the pack has begun
646 // or pre-receive has started, whichever is first.
647 // receive-pack's keepalives count, so indexing and hooks do
648 // not end it. The idle rule waits for the pack because the
649 // client sends nothing while pack-objects counts and
650 // compresses, and receive-pack sends no keepalive then.
651 // Once pre-receive starts only the idle rule applies, so
652 // post-receive is never cut short by the clock.
653 client, clientOut := stdin, stdout
654 var idle <-chan struct{}
655 var arm, unwatch func()
656 stdin, stdout, idle, arm, unwatch = packlimit.Idle(stdin, stdout, idleFor)
657 stdin = &packStart{r: stdin, seen: arm}
658 defer unwatch()
659 started, forget := hookd.AwaitPreReceive(token)
660 defer forget()
661 deadline := time.NewTimer(time.Until(slotAt.Add(receiveFor)))
662 defer deadline.Stop()
663 kill := make(chan struct{})
664 finished := make(chan struct{})
665 defer close(finished)
666 go func() {
667 receiving := deadline.C
668 for {
669 select {
670 case <-finished:
671 return
672 case <-started:
673 arm()
674 started, receiving = nil, nil
675 continue
676 case <-revoked:
677 case <-idle:
678 case <-receiving:
679 }
680 close(kill)
681 // A read blocked on a silent client outlives git;
682 // closing the channel ends it and the stdin copy, so
683 // Transport's Wait returns.
684 for _, c := range []any{client, clientOut} {
685 if c, ok := c.(io.Closer); ok {
686 c.Close()
687 }
688 }
689 return
690 }
691 }()
692 cancel = kill
693 }
694 }
695 if !write {
696 // Pack generation shares one budget with smart HTTP and git://.
697 release, code := takeSlot(packs, "user:"+strconv.FormatInt(user.ID, 10), done, stderr,
698 "the server is busy: it is at its limit of concurrent clones and fetches; try again in a minute")
699 if code != protocol.ExitOK {
700 return code
701 }
702 // Deferred before Transport runs, so it fires after git has
703 // exited and been waited for.
704 defer release()
705 // A client that stops reading would hold its slot for as long
706 // as its channel stays open.
707 client := stdout
708 var stalled <-chan struct{}
709 var unwatch func()
710 stdout, stalled, unwatch = packs.Watch(client)
711 defer unwatch()
712 kill := make(chan struct{})
713 finished := make(chan struct{})
714 defer close(finished)
715 go func() {
716 left := done
717 for {
718 select {
719 case <-finished:
720 return
721 case <-revoked:
722 case <-left:
723 select {
724 case <-stopping:
725 // done closes on a restart too; a clone already
726 // running finishes then. Only a departed client
727 // ends it.
728 left = nil
729 continue
730 default:
731 }
732 case <-stalled:
733 close(kill)
734 // A write blocked on the client's window outlives
735 // git; closing the channel ends it and the stdin copy,
736 // so Transport's Wait returns.
737 if c, ok := client.(io.Closer); ok {
738 c.Close()
739 }
740 return
741 }
742 close(kill)
743 return
744 }
745 }()
746 cancel = kill
747 }
748 if err := gitutil.Transport(service, dir, stdin, stdout, stderr, env, maxPack, keepAlive, cancel); err != nil {
749 return protocol.ExitFailure
750 }
751 return protocol.ExitOK
752}
753
754// takeSlot takes a slot from l for principal, waiting until done closes
755// at most. On a refusal it prints busy, or that the server is going
756// away, and returns a nonzero exit.
757func takeSlot(l *packlimit.Limiter, principal string, done <-chan struct{}, stderr io.Writer, busy string) (release func(), code int) {
758 release, err := l.Acquire(done, principal)
759 if err == nil {
760 return release, protocol.ExitOK
761 }
762 l.Refused("ssh", principal, err)
763 if errors.Is(err, packlimit.ErrBusy) {
764 fmt.Fprintln(stderr, busy)
765 } else {
766 // ErrGone: the client left, or the server is restarting.
767 fmt.Fprintln(stderr, "the server is restarting; try again in a minute")
768 }
769 return nil, protocol.ExitFailure
770}
771
772// pushKeepAlive is receive.keepAlive, in seconds, for a push idle limit
773// of idle: several keepalives fit in one idle period, and never more
774// than five seconds apart.
775func pushKeepAlive(idle time.Duration) int {
776 return max(1, min(5, int(idle/(4*time.Second))))
777}
778
779// packStart calls seen once the pack signature "PACK" has passed
780// through r, which may split it across reads. It can fire early on a
781// command or push option that holds the word; that only starts the idle
782// rule sooner.
783type packStart struct {
784 r io.Reader
785 seen func()
786 tail []byte // up to three bytes carried from the previous read
787 done bool
788}
789
790func (p *packStart) Read(b []byte) (int, error) {
791 n, err := p.r.Read(b)
792 if n > 0 && !p.done {
793 buf := append(p.tail, b[:n]...)
794 if bytes.Contains(buf, []byte("PACK")) {
795 p.done, p.tail = true, nil
796 p.seen()
797 } else {
798 p.tail = append([]byte(nil), buf[max(0, len(buf)-3):]...)
799 }
800 }
801 return n, err
802}