// Package sshd implements the embedded SSH listener: public-key auth against // registered keys, then dispatch to git transport or control commands. package sshd import ( "bytes" "context" "crypto/ed25519" "crypto/rand" "encoding/base64" "encoding/pem" "errors" "fmt" "io" "log/slog" "maps" "net" "os" "path/filepath" "slices" "strconv" "strings" "sync" "sync/atomic" "time" "golang.org/x/crypto/ssh" "gitbay.org/gitbay/internal/config" "gitbay.org/gitbay/internal/control" "gitbay.org/gitbay/internal/gitutil" "gitbay.org/gitbay/internal/hookd" "gitbay.org/gitbay/internal/packlimit" "gitbay.org/gitbay/internal/policy" "gitbay.org/gitbay/internal/protocol" "gitbay.org/gitbay/internal/store" ) type Server struct { cfg config.Config st *store.Store packs *packlimit.Limiter pushes *packlimit.Limiter sshCfg *ssh.ServerConfig authLimiter *rateLimiter sessions sync.WaitGroup // accepted connections still being served mu sync.Mutex conns map[*conn]struct{} stopping chan struct{} // closed by Stop stopOnce sync.Once } // conn is one accepted connection and how many sessions it is running. // A CLI's shared connection sits idle between commands; on shutdown an // idle connection is closed at once and only a session mid-command is // waited for (#141). type conn struct { net net.Conn active atomic.Int32 // keyID and userID are the key that authenticated the connection and // its account: 0 before the handshake and for an unregistered key. // Guarded by Server.mu. keyID, userID int64 revoked chan struct{} // closed by cut cutOnce sync.Once } // cut ends the connection because its key was revoked: a git transport // on it is killed, and every other command loses its channel. func (c *conn) cut() { c.cutOnce.Do(func() { close(c.revoked) }) c.net.Close() } func New(cfg config.Config, st *store.Store, packs, pushes *packlimit.Limiter) (*Server, error) { 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{})} sc := &ssh.ServerConfig{ PublicKeyCallback: s.authenticate, ServerVersion: "SSH-2.0-gitbayd", } signers, err := loadHostKeys(cfg) if err != nil { return nil, err } for _, sg := range signers { sc.AddHostKey(sg) } s.sshCfg = sc st.OnRevoke(s.revoke) return s, nil } // loadHostKeys loads the configured host keys, or generates an ed25519 key // under server.root/ssh/ when none are configured. func loadHostKeys(cfg config.Config) ([]ssh.Signer, error) { paths := cfg.SSH.HostKeys if len(paths) == 0 { p := filepath.Join(cfg.Server.Root, "ssh", "host_ed25519") if _, err := os.Stat(p); errors.Is(err, os.ErrNotExist) { if err := generateHostKey(p); err != nil { return nil, fmt.Errorf("generating host key: %w", err) } slog.Info("generated ssh host key", "path", p) } paths = []string{p} } var signers []ssh.Signer for _, p := range paths { raw, err := os.ReadFile(p) if err != nil { return nil, fmt.Errorf("host key %s: %w", p, err) } sg, err := ssh.ParsePrivateKey(raw) if err != nil { return nil, fmt.Errorf("host key %s: %w", p, err) } signers = append(signers, sg) } return signers, nil } func generateHostKey(path string) error { if err := os.MkdirAll(filepath.Dir(path), 0o700); err != nil { return err } _, priv, err := ed25519.GenerateKey(rand.Reader) if err != nil { return err } block, err := ssh.MarshalPrivateKey(priv, "") if err != nil { return err } return os.WriteFile(path, pem.EncodeToMemory(block), 0o600) } // authenticate resolves the presented key to a registered account. The SSH // username is ignored; identity comes from the key alone. When registration // is open or invite-based, unknown keys are admitted to run exactly one // command: register. func (s *Server) authenticate(meta ssh.ConnMetadata, pub ssh.PublicKey) (*ssh.Permissions, error) { ip := remoteIP(meta.RemoteAddr()) if !s.authLimiter.allow(ip) { // One audit entry per throttled window, not per rejected attempt. if s.authLimiter.firstThrottle(ip) { s.st.Audit(0, "auth.throttled", map[string]any{"ip": ip, "rate": s.cfg.Limits.SSHAuthRate}) } return nil, fmt.Errorf("too many authentication attempts; try again shortly") } fp := ssh.FingerprintSHA256(pub) key, err := s.st.SSHKeyByFingerprint(fp) if err != nil && !errors.Is(err, store.ErrNotFound) { // The store, not the key, failed. Neither a failure against the // limiter nor "unknown key": a busy database during a restart // would otherwise lock every client out for a minute. slog.Error("ssh auth: key lookup", "err", err) return nil, fmt.Errorf("authentication temporarily unavailable") } if err != nil { if s.cfg.Registration.Mode != "closed" { return &ssh.Permissions{Extensions: map[string]string{ "anon-key": base64.StdEncoding.EncodeToString(pub.Marshal()), }}, nil } s.authLimiter.fail(ip) s.st.Audit(0, "auth.failed", map[string]any{"ip": ip, "fingerprint": fp}) return nil, fmt.Errorf("unknown key %s", fp) } if key.Expired(time.Now()) { s.authLimiter.fail(ip) s.st.Audit(key.UserID, "auth.expired", map[string]any{"ip": ip, "fingerprint": fp}) return nil, fmt.Errorf("key %s has expired", fp) } s.authLimiter.success(ip) return &ssh.Permissions{Extensions: map[string]string{ "user-id": strconv.FormatInt(key.UserID, 10), "key-id": strconv.FormatInt(key.ID, 10), }}, nil } // Serve accepts connections on ln until it is closed. func (s *Server) Serve(ln net.Listener) error { served := make(chan struct{}) defer close(served) go s.sweep(served) for { nc, err := ln.Accept() if err != nil { return err } c := &conn{net: nc, revoked: make(chan struct{})} s.mu.Lock() s.conns[c] = struct{}{} s.mu.Unlock() s.sessions.Add(1) go func() { defer s.sessions.Done() defer func() { s.mu.Lock() delete(s.conns, c) s.mu.Unlock() }() s.handleConn(c) }() } } // revoke closes the connections opened by the keys r names. func (s *Server) revoke(r store.Revoked) { var cut []*conn s.mu.Lock() for c := range s.conns { if c.keyID == 0 { continue } if (r.UserID != 0 && c.userID == r.UserID) || slices.Contains(r.KeyIDs, c.keyID) { cut = append(cut, c) } } s.mu.Unlock() for _, c := range cut { c.cut() } } // sweepInterval bounds how long a revocation this process was not told // about (gitbayd admin on the host) leaves a connection open. const sweepInterval = 15 * time.Second func (s *Server) sweep(served <-chan struct{}) { t := time.NewTicker(sweepInterval) defer t.Stop() for { select { case <-t.C: s.sweepOnce() case <-served: return case <-s.stopping: return } } } // sweepOnce cuts every connection whose key is no longer live. Only // connections whose key was asked about are judged: one that // authenticated while the query ran waits for the next sweep. func (s *Server) sweepOnce() { asked := map[int64]bool{} s.mu.Lock() for c := range s.conns { if c.keyID != 0 { asked[c.keyID] = true } } s.mu.Unlock() if len(asked) == 0 { return } live, err := s.st.LiveSSHKeys(slices.Collect(maps.Keys(asked))) if err != nil { slog.Error("ssh sweep: key lookup", "err", err) return } var cut []*conn s.mu.Lock() for c := range s.conns { if asked[c.keyID] && !live[c.keyID] { cut = append(cut, c) } } s.mu.Unlock() for _, c := range cut { c.cut() } } // Stop ends the commands that run until something happens (build log // --follow), so a shutdown drain waits only for work that finishes. It // does not close connections; Shutdown does. func (s *Server) Stop() { s.stopOnce.Do(func() { close(s.stopping) }) } // Shutdown closes every idle connection, then waits for the ones with a // session running, or for ctx. The caller closes the listener first; a // push in flight completes rather than being cut mid-pack. func (s *Server) Shutdown(ctx context.Context) error { s.Stop() s.mu.Lock() for c := range s.conns { if c.active.Load() == 0 { c.net.Close() } } s.mu.Unlock() done := make(chan struct{}) go func() { s.sessions.Wait() close(done) }() select { case <-done: return nil case <-ctx.Done(): return ctx.Err() } } func (s *Server) handleConn(c *conn) { defer c.net.Close() sconn, chans, reqs, err := ssh.NewServerConn(c.net, s.sshCfg) if err != nil { return } defer sconn.Close() ext := sconn.Permissions.Extensions s.mu.Lock() c.keyID, _ = strconv.ParseInt(ext["key-id"], 10, 64) c.userID, _ = strconv.ParseInt(ext["user-id"], 10, 64) s.mu.Unlock() go ssh.DiscardRequests(reqs) for newCh := range chans { if newCh.ChannelType() != "session" { newCh.Reject(ssh.UnknownChannelType, "only session channels are supported") continue } ch, chReqs, err := newCh.Accept() if err != nil { continue } c.active.Add(1) go func() { defer c.active.Add(-1) s.handleSession(c, sconn, ch, chReqs) }() } } func (s *Server) handleSession(c *conn, sconn *ssh.ServerConn, ch ssh.Channel, reqs <-chan *ssh.Request) { defer ch.Close() var term control.Term for req := range reqs { switch req.Type { case "exec": var payload struct{ Command string } if err := ssh.Unmarshal(req.Payload, &payload); err != nil { req.Reply(false, nil) continue } req.Reply(true, nil) // x/crypto closes reqs when the client closes the channel. That // is how a follow learns nobody is reading: the CLI's shared // connection outlives a Ctrl-C, the channel does not. Stop // ends it too, for a restart. closed := make(chan struct{}) go func() { for r := range reqs { r.Reply(false, nil) } close(closed) }() done := make(chan struct{}) go func() { select { case <-closed: case <-s.stopping: } close(done) }() code := s.runExec(c, sconn, ch, term, payload.Command, done) sendExit(ch, code) return case "shell": req.Reply(true, nil) fmt.Fprintf(ch, "gitbay control plane: interactive shells are not available.\nTry: ssh %s help\n", s.cfg.Server.SiteURL) sendExit(ch, protocol.ExitUsage) return case "env": var kv struct{ Name, Value string } if ssh.Unmarshal(req.Payload, &kv) == nil && kv.Name == "GITBAY_TERM" { term = control.ParseTerm(kv.Value) } req.Reply(true, nil) case "pty-req": // Harmless; accept and ignore. req.Reply(true, nil) default: req.Reply(false, nil) } } } func sendExit(ch ssh.Channel, code int) { var msg = struct{ Status uint32 }{uint32(code)} ch.SendRequest("exit-status", false, ssh.Marshal(&msg)) } func (s *Server) runExec(c *conn, sconn *ssh.ServerConn, ch ssh.Channel, term control.Term, cmdline string, done <-chan struct{}) int { ext := sconn.Permissions.Extensions if blob := ext["anon-key"]; blob != "" { return s.runAnonymous(ch, blob, cmdline) } userID, _ := strconv.ParseInt(ext["user-id"], 10, 64) keyID, _ := strconv.ParseInt(ext["key-id"], 10, 64) // A connection outlives its commands, so the key is read again for // each one: what it may do is what it may do now (#256). key, err := s.st.SSHKeyByID(keyID) if errors.Is(err, store.ErrNotFound) || (err == nil && key.UserID != userID) { fmt.Fprintln(ch.Stderr(), "this key is no longer registered") return protocol.ExitDenied } if err != nil { slog.Error("ssh exec: key lookup", "err", err) fmt.Fprintln(ch.Stderr(), "authentication temporarily unavailable") return protocol.ExitFailure } if key.Expired(time.Now()) { fmt.Fprintln(ch.Stderr(), "this key has expired; remove it and add a new one") return protocol.ExitDenied } user, err := s.st.UserByID(userID) if err != nil { fmt.Fprintln(ch.Stderr(), "account no longer exists") return protocol.ExitDenied } _ = s.st.TouchSSHKey(keyID) return Exec(s.cfg, s.st, s.packs, s.pushes, user, key, term, cmdline, ch, ch, ch.Stderr(), done, s.stopping, c.revoked) } // runAnonymous handles a session from an unregistered key: the register // command and nothing else. func (s *Server) runAnonymous(ch ssh.Channel, keyB64, cmdline string) int { raw, err := base64.StdEncoding.DecodeString(keyB64) if err != nil { return protocol.ExitFailure } pub, err := ssh.ParsePublicKey(raw) if err != nil { return protocol.ExitFailure } argv, err := protocol.Tokenize(cmdline) if err != nil { fmt.Fprintf(ch.Stderr(), "cannot parse command: %v\n", err) return protocol.ExitUsage } if len(argv) == 0 || argv[0] != "register" { host := s.cfg.SiteHost() fp := ssh.FingerprintSHA256(pub) flag := map[string]string{"open": "--email
", "invite": "--invite "}[s.cfg.Registration.Mode] fmt.Fprintf(ch.Stderr(), "this key (%s) is not registered on %s.\n"+ "already have an account? add it at %s/settings#keys\n"+ "new here? ssh git@%s register --username %s\n", fp, host, strings.TrimSuffix(s.cfg.Server.SiteURL, "/"), host, flag) return protocol.ExitDenied } return control.RunRegister(s.cfg, s.st, pub, argv, ch, ch.Stderr()) } // Exec runs one SSH exec command line for an authenticated key. It is the // single dispatch path shared by the embedded listener and the system-sshd // forced command (gitbayd shell). Closing revoked kills a git transport. // packs bounds clones and fetches, and repo download; pushes bounds // receive-pack. A nil limiter is no limit. func Exec(cfg config.Config, st *store.Store, packs, pushes *packlimit.Limiter, user store.User, key store.SSHKey, term control.Term, cmdline string, stdin io.Reader, stdout, stderr io.Writer, done, stopping, revoked <-chan struct{}) int { // A person signing in with a full-scope key cancels a scheduled // deletion; automation on narrower keys is refused and cannot. if user.Disabled && user.DeleteAfter != "" && key.Scope == "full" && control.CancelScheduledDeletion(st, &user, "ssh") { fmt.Fprintln(stderr, "deletion of your account was cancelled") } if user.Disabled { if user.DeleteAfter != "" { fmt.Fprintln(stderr, control.ScheduledRefusal(user)) return protocol.ExitDenied } fmt.Fprintln(stderr, "this account is disabled; contact the instance admin") return protocol.ExitDenied } argv, err := protocol.Tokenize(cmdline) if err != nil { fmt.Fprintf(stderr, "cannot parse command: %v\n", err) return protocol.ExitUsage } if len(argv) > 0 { switch argv[0] { case "git-upload-pack", "git-receive-pack", "git-upload-archive": code := protocol.ExitDenied if user.Pending { fmt.Fprintln(stderr, "your account is not active yet: verify your email first") } else { code = runGit(cfg, st, packs, pushes, user, key, argv, stdin, stdout, stderr, done, stopping, revoked) } // A refused push is a refused write, audited like one. runGit // refuses only with the path as the one argument, so argv[1:] // holds no value beyond the target. if argv[0] == "git-receive-pack" && (code == protocol.ExitDenied || code == protocol.ExitNotFound) { control.AuditRefused(st, user.ID, "refused git-receive-pack", map[string]any{"argv": argv[1:], "source": key.Fingerprint, "exit": code}) } return code case "git-lfs-authenticate": // Part of the git transport, not the control plane: usable by // git-scoped and deploy keys, with the transports' access rules. if user.Pending { fmt.Fprintln(stderr, "your account is not active yet: verify your email first") return protocol.ExitDenied } return runLFSAuthenticate(cfg, st, user, key, argv, stdout, stderr) } } ctx := &control.Ctx{ User: user, Scope: key.Scope, Source: key.Fingerprint, Term: term, Store: st, Cfg: cfg, Stdin: stdin, Stdout: stdout, Stderr: stderr, Done: done, Stopping: stopping, Expires: key.ExpiresAt, Packs: packs, } return control.Dispatch(ctx, argv) } // runGit streams a git transport service after access checks. func runGit(cfg config.Config, st *store.Store, packs, pushes *packlimit.Limiter, user store.User, key store.SSHKey, argv []string, stdin io.Reader, stdout, stderr io.Writer, done, stopping, revoked <-chan struct{}) int { service, scope := argv[0], key.Scope if len(argv) != 2 { fmt.Fprintf(stderr, "usage: %s \n", service) return protocol.ExitUsage } write := service == "git-receive-pack" cancel, keepAlive := revoked, 0 repo, err := st.RepoByPath(argv[1]) if err != nil { fmt.Fprintln(stderr, "repository not found") return protocol.ExitNotFound } if policy.IsDeployScope(scope) { // A deploy key authorizes by its binding alone: one repository, // its mode, nothing inherited from whoever registered it. Any // mismatch reads as nonexistence, same as the access rules. if !policy.DeployScopeAllows(scope, repo.ID, write) { fmt.Fprintln(stderr, "repository not found") return protocol.ExitNotFound } } else { grant, err := st.AccessRole(repo.ID, user.ID) if err != nil { fmt.Fprintln(stderr, "internal error") return protocol.ExitFailure } if !policy.CanRead(user, repo, grant) { // Same answer as nonexistence: private repos must not be enumerable. fmt.Fprintln(stderr, "repository not found") return protocol.ExitNotFound } if !policy.ScopeAllowsGit(scope, repo.Path(), write) { fmt.Fprintf(stderr, "this key's scope (%s) does not allow %s on %s\n", scope, service, repo.Path()) return protocol.ExitDenied } if write && !policy.CanWrite(user, repo, grant) { fmt.Fprintf(stderr, "write access to %s denied\n", repo.Path()) return protocol.ExitDenied } } if write && repo.Settings.Archived { fmt.Fprintf(stderr, "%s is archived and read-only\n", repo.Path()) return protocol.ExitDenied } if write { if mirrored, err := st.PullMirrored(repo.ID); err == nil && mirrored { fmt.Fprintf(stderr, "%s is a pull mirror: its refs come from the upstream; push there instead\n", repo.Path()) return protocol.ExitDenied } } dir := control.RepoDir(cfg.Server.Root, repo.OwnerName, repo.Name) env := []string{ hookd.EnvSocket + "=" + hookd.SocketPath(cfg.Server.Root), hookd.EnvRepoID + "=" + strconv.FormatInt(repo.ID, 10), hookd.EnvUserID + "=" + strconv.FormatInt(user.ID, 10), hookd.EnvScope + "=" + scope, } // A storage quota on the owner rides the same mechanism as the pack // cap: the pack may be no larger than what the owner has left. maxPack := cfg.Limits.MaxPackBytes if write { if limit := control.ByteLimit(st, control.QuotaConfig(cfg), repo.OwnerKind, repo.OwnerID); limit > 0 { used := control.OwnedBytes(st, cfg.Server.Root, repo.OwnerKind, repo.OwnerID) left := limit - used if left <= 0 { 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) return protocol.ExitDenied } if maxPack == 0 || left < maxPack { maxPack = left } } } if write { // Pushes have their own budget, so a clone storm cannot starve // them or the reverse. The slot covers receive-pack and both // hooks: git waits for post-receive (RefsUpdated) before it // exits, and that work is the push's cost. A deploy key is its // own principal, not the account that registered it. principal := "user:" + strconv.FormatInt(user.ID, 10) if policy.IsDeployScope(scope) { principal = "key:" + strconv.FormatInt(key.ID, 10) } release, code := takeSlot(pushes, principal, done, stderr, "the server is busy: it is at its limit of concurrent pushes; try again in a minute") if code != protocol.ExitOK { return code } slotAt := time.Now() // Deferred before Transport runs, so it fires after receive-pack // has exited, on every path: the client hanging up ends its // stdin and receive-pack with it, and a revoked key kills it. defer release() // hookd answers only a hook that names this receive-pack. token, err := st.CreatePushToken(repo.ID, user.ID, scope) if err != nil { fmt.Fprintln(stderr, "internal error") return protocol.ExitFailure } defer st.DeletePushToken(token) env = append(env, hookd.EnvToken+"="+token) idleFor, receiveFor := cfg.Limits.PushTimeouts() keepAlive = pushKeepAlive(idleFor) if pushes != nil { // A client that holds its slot while sending nothing, or // trickles its pack, is cut: when pre-receive has not // started receiveFor after the slot was taken, or after // idleFor with no byte either way once the pack has begun // or pre-receive has started, whichever is first. // receive-pack's keepalives count, so indexing and hooks do // not end it. The idle rule waits for the pack because the // client sends nothing while pack-objects counts and // compresses, and receive-pack sends no keepalive then. // Once pre-receive starts only the idle rule applies, so // post-receive is never cut short by the clock. client, clientOut := stdin, stdout var idle <-chan struct{} var arm, unwatch func() stdin, stdout, idle, arm, unwatch = packlimit.Idle(stdin, stdout, idleFor) stdin = &packStart{r: stdin, seen: arm} defer unwatch() started, forget := hookd.AwaitPreReceive(token) defer forget() deadline := time.NewTimer(time.Until(slotAt.Add(receiveFor))) defer deadline.Stop() kill := make(chan struct{}) finished := make(chan struct{}) defer close(finished) go func() { receiving := deadline.C for { select { case <-finished: return case <-started: arm() started, receiving = nil, nil continue case <-revoked: case <-idle: case <-receiving: } close(kill) // A read blocked on a silent client outlives git; // closing the channel ends it and the stdin copy, so // Transport's Wait returns. for _, c := range []any{client, clientOut} { if c, ok := c.(io.Closer); ok { c.Close() } } return } }() cancel = kill } } if !write { // Pack generation shares one budget with smart HTTP and git://. release, code := takeSlot(packs, "user:"+strconv.FormatInt(user.ID, 10), done, stderr, "the server is busy: it is at its limit of concurrent clones and fetches; try again in a minute") if code != protocol.ExitOK { return code } // Deferred before Transport runs, so it fires after git has // exited and been waited for. defer release() // A client that stops reading would hold its slot for as long // as its channel stays open. client := stdout var stalled <-chan struct{} var unwatch func() stdout, stalled, unwatch = packs.Watch(client) defer unwatch() kill := make(chan struct{}) finished := make(chan struct{}) defer close(finished) go func() { left := done for { select { case <-finished: return case <-revoked: case <-left: select { case <-stopping: // done closes on a restart too; a clone already // running finishes then. Only a departed client // ends it. left = nil continue default: } case <-stalled: close(kill) // A write blocked on the client's window outlives // git; closing the channel ends it and the stdin copy, // so Transport's Wait returns. if c, ok := client.(io.Closer); ok { c.Close() } return } close(kill) return } }() cancel = kill } if err := gitutil.Transport(service, dir, stdin, stdout, stderr, env, maxPack, keepAlive, cancel); err != nil { return protocol.ExitFailure } return protocol.ExitOK } // takeSlot takes a slot from l for principal, waiting until done closes // at most. On a refusal it prints busy, or that the server is going // away, and returns a nonzero exit. func takeSlot(l *packlimit.Limiter, principal string, done <-chan struct{}, stderr io.Writer, busy string) (release func(), code int) { release, err := l.Acquire(done, principal) if err == nil { return release, protocol.ExitOK } l.Refused("ssh", principal, err) if errors.Is(err, packlimit.ErrBusy) { fmt.Fprintln(stderr, busy) } else { // ErrGone: the client left, or the server is restarting. fmt.Fprintln(stderr, "the server is restarting; try again in a minute") } return nil, protocol.ExitFailure } // pushKeepAlive is receive.keepAlive, in seconds, for a push idle limit // of idle: several keepalives fit in one idle period, and never more // than five seconds apart. func pushKeepAlive(idle time.Duration) int { return max(1, min(5, int(idle/(4*time.Second)))) } // packStart calls seen once the pack signature "PACK" has passed // through r, which may split it across reads. It can fire early on a // command or push option that holds the word; that only starts the idle // rule sooner. type packStart struct { r io.Reader seen func() tail []byte // up to three bytes carried from the previous read done bool } func (p *packStart) Read(b []byte) (int, error) { n, err := p.r.Read(b) if n > 0 && !p.done { buf := append(p.tail, b[:n]...) if bytes.Contains(buf, []byte("PACK")) { p.done, p.tail = true, nil p.seen() } else { p.tail = append([]byte(nil), buf[max(0, len(buf)-3):]...) } } return n, err }