Commit be2c840ec8
Verified · cmc
Layout: unified · split
cmd/gitbayd/main.go +1 −1
| @@ -232,7 +232,7 @@ func serveCmd() *cobra.Command { | ||
| 232 | 232 | var sshSrv *sshd.Server |
| 233 | 233 | var sshLn, gitLn net.Listener |
| 234 | 234 | if cfg.SSH.Mode == "embedded" { |
| 235 | srv, err := sshd.New(cfg, st) | |
| 235 | srv, err := sshd.New(cfg, st, nil) | |
| 236 | 236 | if err != nil { |
| 237 | 237 | return err |
| 238 | 238 | } |
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 |
internal/gitutil/gitutil.go +5
| @@ -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; sshd 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) |
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 stallDeadline. | |
| 196 | func TestCloneKilledWhenClientStopsReading(t *testing.T) { | |
| 197 | old := stallDeadline | |
| 198 | stallDeadline = 200 * time.Millisecond | |
| 199 | t.Cleanup(func() { 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 +99 −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:] |
| @@ -514,9 +516,28 @@ func Exec(cfg config.Config, st *store.Store, user store.User, key store.SSHKey, | ||
| 514 | 516 | return control.Dispatch(ctx, argv) |
| 515 | 517 | } |
| 516 | 518 | |
| 519 | // stallDeadline is how long a limited transport may go without writing | |
| 520 | // to its client before it is killed. upload-pack sends a keepalive | |
| 521 | // every five seconds while it prepares a pack. | |
| 522 | var stallDeadline = 2 * time.Minute | |
| 523 | ||
| 524 | // progressWriter records when a write to the client last completed. | |
| 525 | type progressWriter struct { | |
| 526 | w io.Writer | |
| 527 | last atomic.Int64 // unix nanoseconds | |
| 528 | } | |
| 529 | ||
| 530 | func (p *progressWriter) Write(b []byte) (int, error) { | |
| 531 | n, err := p.w.Write(b) | |
| 532 | if n > 0 { | |
| 533 | p.last.Store(time.Now().UnixNano()) | |
| 534 | } | |
| 535 | return n, err | |
| 536 | } | |
| 537 | ||
| 517 | 538 | // 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 { | |
| 539 | func runGit(cfg config.Config, st *store.Store, packs *packlimit.Limiter, user store.User, scope string, argv []string, | |
| 540 | stdin io.Reader, stdout, stderr io.Writer, done, stopping, revoked <-chan struct{}) int { | |
| 520 | 541 | service := argv[0] |
| 521 | 542 | if len(argv) != 2 { |
| 522 | 543 | fmt.Fprintf(stderr, "usage: %s <path>\n", service) |
| @@ -601,7 +622,77 @@ func runGit(cfg config.Config, st *store.Store, user store.User, scope string, a | ||
| 601 | 622 | defer st.DeletePushToken(token) |
| 602 | 623 | env = append(env, hookd.EnvToken+"="+token) |
| 603 | 624 | } |
| 604 | if err := gitutil.Transport(service, dir, stdin, stdout, stderr, env, maxPack, revoked); err != nil { | |
| 625 | cancel := revoked | |
| 626 | if !write { | |
| 627 | // Pack generation shares one budget with smart HTTP and git://. | |
| 628 | // receive-pack stays outside it: its post-receive runs after the | |
| 629 | // client has its report, and must not be queued or killed. | |
| 630 | release, err := packs.Acquire(done, "user:"+strconv.FormatInt(user.ID, 10)) | |
| 631 | if errors.Is(err, packlimit.ErrBusy) { | |
| 632 | fmt.Fprintln(stderr, "the server is busy: it is at its limit of concurrent clones and fetches; try again in a minute") | |
| 633 | return protocol.ExitFailure | |
| 634 | } | |
| 635 | if err != nil { | |
| 636 | // ErrGone: the client left, or the server is restarting. | |
| 637 | fmt.Fprintln(stderr, "the server is restarting; try again in a minute") | |
| 638 | return protocol.ExitFailure | |
| 639 | } | |
| 640 | // Deferred before Transport runs, so it fires after git has | |
| 641 | // exited and been waited for. | |
| 642 | defer release() | |
| 643 | // A client that stops reading would hold its slot for as long | |
| 644 | // as its channel stays open. With a limit in force, a transport | |
| 645 | // that writes nothing for stallDeadline is killed. | |
| 646 | var pw *progressWriter | |
| 647 | var tick <-chan time.Time | |
| 648 | if packs != nil { | |
| 649 | pw = &progressWriter{w: stdout} | |
| 650 | pw.last.Store(time.Now().UnixNano()) | |
| 651 | stdout = pw | |
| 652 | t := time.NewTicker(stallDeadline / 4) | |
| 653 | defer t.Stop() | |
| 654 | tick = t.C | |
| 655 | } | |
| 656 | kill := make(chan struct{}) | |
| 657 | finished := make(chan struct{}) | |
| 658 | defer close(finished) | |
| 659 | go func() { | |
| 660 | left := done | |
| 661 | for { | |
| 662 | select { | |
| 663 | case <-finished: | |
| 664 | return | |
| 665 | case <-revoked: | |
| 666 | case <-left: | |
| 667 | select { | |
| 668 | case <-stopping: | |
| 669 | // done closes on a restart too; a clone already | |
| 670 | // running finishes then. Only a departed client | |
| 671 | // ends it. | |
| 672 | left = nil | |
| 673 | continue | |
| 674 | default: | |
| 675 | } | |
| 676 | case <-tick: | |
| 677 | if time.Since(time.Unix(0, pw.last.Load())) < stallDeadline { | |
| 678 | continue | |
| 679 | } | |
| 680 | close(kill) | |
| 681 | // A write blocked on the client's window outlives | |
| 682 | // git; closing the channel ends it and the stdin copy, | |
| 683 | // so Transport's Wait returns. | |
| 684 | if c, ok := pw.w.(io.Closer); ok { | |
| 685 | c.Close() | |
| 686 | } | |
| 687 | return | |
| 688 | } | |
| 689 | close(kill) | |
| 690 | return | |
| 691 | } | |
| 692 | }() | |
| 693 | cancel = kill | |
| 694 | } | |
| 695 | if err := gitutil.Transport(service, dir, stdin, stdout, stderr, env, maxPack, cancel); err != nil { | |
| 605 | 696 | return protocol.ExitFailure |
| 606 | 697 | } |
| 607 | 698 | return protocol.ExitOK |
internal/sshd/sshd_test.go +1 −1
| @@ -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 | } |