internal/sshd/sshd.go

e6cd75b5f28bacf51620bb531320c30fd4e66bfd
gitbay/internal/sshd/sshd.go history · blame · raw

802 lines · 25331 bytes

25 symbols in this file
  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}