diff --git a/.cspell.json b/.cspell.json index 870ee86d048..35c2d6e5ad9 100644 --- a/.cspell.json +++ b/.cspell.json @@ -160,6 +160,7 @@ "opensource", "opentype", "Pacman", + "perr", "picus", "Pinia", "pkce", @@ -230,6 +231,7 @@ "unplugin", "unsanitize", "Upsert", + "upstack", "urfave", "usecase", "useragent", @@ -250,6 +252,7 @@ "woodpeckerci", "WORKDIR", "Wrapf", + "writedeadline", "x-enum-varnames", "xlink", "xlog", diff --git a/cmd/server/server.go b/cmd/server/server.go index 240b918941d..a52cbd1504a 100644 --- a/cmd/server/server.go +++ b/cmd/server/server.go @@ -50,8 +50,33 @@ import ( const ( shutdownTimeout = time.Second * 5 + + readHeaderTimeout = 10 * time.Second + writeTimeout = 60 * time.Second + idleTimeout = 120 * time.Second ) +// setupServerTimeouts applies conservative timeouts to an http.Server so a +// response write blocked on a zero-window peer, or an abandoned keepalive +// connection, is reclaimed instead of orphaning the connection forever. +// +// ReadTimeout is intentionally left zero: request bodies here are small +// (webhooks, API JSON — agent log upload rides gRPC on :9000, not this +// server), so an unbounded body read is low exposure. The slow-loris variants +// are still bounded: a slow HEADER is capped by ReadHeaderTimeout, and a slow +// BODY is capped by WriteTimeout — net/http arms the write deadline once the +// headers are read, so a handler that has begun responding cannot be held open +// indefinitely by a trickled request body. A future body-streaming handler +// that re-arms its own write deadline (as the SSE handlers do) would escape +// that WriteTimeout bound and must add its own body/read deadline. SSE +// handlers override WriteTimeout per-response via http.ResponseController (see +// server/api/stream.go). +func setupServerTimeouts(s *http.Server) { + s.ReadHeaderTimeout = readHeaderTimeout + s.WriteTimeout = writeTimeout + s.IdleTimeout = idleTimeout +} + var ( stopServerFunc context.CancelCauseFunc = func(error) {} shutdownCancelFunc context.CancelFunc = func() {} @@ -192,6 +217,7 @@ func run(ctx context.Context, c *cli.Command) error { NextProtos: []string{"h2", "http/1.1"}, }, } + setupServerTimeouts(tlsServer) go func() { <-ctx.Done() @@ -231,6 +257,7 @@ func run(ctx context.Context, c *cli.Command) error { Addr: server.Config.Server.Port, Handler: http.HandlerFunc(redirect), } + setupServerTimeouts(redirectServer) go func() { <-ctx.Done() log.Info().Msg("shutdown redirect server ...") @@ -278,6 +305,7 @@ func run(ctx context.Context, c *cli.Command) error { httpServer := &http.Server{ Handler: handler, } + setupServerTimeouts(httpServer) go func() { <-ctx.Done() @@ -307,6 +335,7 @@ func run(ctx context.Context, c *cli.Command) error { Addr: metricsServerAddr, Handler: metricsRouter, } + setupServerTimeouts(metricsServer) go func() { <-ctx.Done() diff --git a/cmd/server/server_test.go b/cmd/server/server_test.go index 44948181470..d28f642802b 100644 --- a/cmd/server/server_test.go +++ b/cmd/server/server_test.go @@ -4,7 +4,7 @@ // you may not use this file except in compliance with the License. // You may obtain a copy of the License at // -// http://www.apache.org/licenses/LICENSE-2.0 +// http://www.apache.org/licenses/LICENSE-2.0 // // Unless required by applicable law or agreed to in writing, software // distributed under the License is distributed on an "AS IS" BASIS, @@ -15,52 +15,302 @@ package main import ( + "bufio" + "go/ast" + "go/parser" + "go/token" + "io" + "net" + "net/http" "os" + "strings" "testing" + "time" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) -func TestSetUnixSocketPermission(t *testing.T) { - tests := []struct { - name string - permission string - want os.FileMode - wantErr bool - }{ - { - name: "owner and group read write", - permission: "660", - want: 0o660, - }, - { - name: "world read write", - permission: "666", - want: 0o666, - }, - { - name: "rejects non-octal digits", - permission: "999", - wantErr: true, - }, +// TestServerTimeoutsWedgeKillBlockedWrite pins the fix: a response write +// blocked on a peer that stops reading (a zero-window wedge) must be +// reclaimed by the server-wide WriteTimeout instead of orphaning the +// connection forever. With setupServerTimeouts applied the handler's Write +// returns os.ErrDeadlineExceeded and the connection is closed within +// WriteTimeout + slack. The applied WriteTimeout is shortened below to keep +// the suite fast, so this test pins the reclaim mechanism; that the helper +// actually sets the field is pinned by TestSetupServerTimeoutsFieldValues. +func TestServerTimeoutsWedgeKillBlockedWrite(t *testing.T) { + // Short test WriteTimeout mirrors the helper's field shape (not the prod + // 60s) to keep the suite fast. The mechanism under test is identical. + const testWriteTimeout = 500 * time.Millisecond + const slack = 3 * time.Second + + writeErrCh := make(chan error, 1) + handler := http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Content-Type", "application/octet-stream") + // A multi-MB body: once the client stops reading and the send buffer + // fills, this Write blocks and the write deadline must reclaim it. + payload := make([]byte, 32<<20) + _, err := w.Write(payload) + writeErrCh <- err + }) + + ln, err := net.Listen("tcp", "127.0.0.1:0") + require.NoError(t, err) + + srv := &http.Server{Handler: handler} + setupServerTimeouts(srv) // GREEN: the fix under test. + // Assert before overwriting: shortening the field for speed would otherwise + // mask a helper that never set it, leaving this test green against exactly + // the regression it exists to catch. + require.NotZero(t, srv.WriteTimeout, "setupServerTimeouts must set WriteTimeout") + srv.WriteTimeout = testWriteTimeout // shorten the applied field for speed. + go func() { _ = srv.Serve(ln) }() + t.Cleanup(func() { _ = srv.Close() }) + + conn, err := net.Dial("tcp", ln.Addr().String()) + require.NoError(t, err) + t.Cleanup(func() { _ = conn.Close() }) + + _, err = conn.Write([]byte("GET / HTTP/1.1\r\nHost: test\r\n\r\n")) + require.NoError(t, err) + + // Read only the first bytes, then stop reading (do NOT close) so the + // server-side Write blocks on a full send buffer and the write deadline, + // anchored at request-header read, fires. + firstByte := make([]byte, 1024) + _, err = conn.Read(firstByte) + require.NoError(t, err) + + select { + case werr := <-writeErrCh: + require.Error(t, werr, "handler Write unexpectedly succeeded") + assert.ErrorIs(t, werr, os.ErrDeadlineExceeded, + "wedged handler Write should fail with a deadline error, got: %v", werr) + case <-time.After(testWriteTimeout + slack): + t.Fatalf("handler Write did not return within WriteTimeout+slack (%s); "+ + "the wedged write was never bounded", testWriteTimeout+slack) } - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - path := t.TempDir() + "/socket" - require.NoError(t, os.WriteFile(path, []byte(""), 0o600)) + // The server closes the wedged connection: a subsequent read drains the + // buffered bytes then observes the close (EOF/reset) within slack. A + // deadline-exceeded here would mean the conn was never reclaimed. + require.NoError(t, conn.SetReadDeadline(time.Now().Add(slack))) + drain := make([]byte, 4096) + for { + _, rerr := conn.Read(drain) + if rerr != nil { + assert.NotErrorIs(t, rerr, os.ErrDeadlineExceeded, + "server did not close the wedged connection within slack") + break + } + } +} + +// TestServerTimeoutsIdleOrphanKill pins the idle-connection half of the fix: a +// client that completes a request, stops reading, and holds the keepalive +// connection open without issuing another request must have that connection +// reclaimed by IdleTimeout. Against a bare no-timeout server the connection is +// never closed and the read below trips its deadline (red). +func TestServerTimeoutsIdleOrphanKill(t *testing.T) { + const testIdleTimeout = 500 * time.Millisecond + const slack = 3 * time.Second + + handler := http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Content-Type", "text/plain") + // Small response: fits the kernel send buffer, handler returns cleanly, + // leaving an idle keepalive connection. + _, _ = io.WriteString(w, "ok") + }) + + ln, err := net.Listen("tcp", "127.0.0.1:0") + require.NoError(t, err) + + srv := &http.Server{Handler: handler} + setupServerTimeouts(srv) // GREEN: the fix under test. + // Assert before overwriting — see the note in the wedge-kill test. + require.NotZero(t, srv.IdleTimeout, "setupServerTimeouts must set IdleTimeout") + srv.IdleTimeout = testIdleTimeout // shorten the applied field for speed. + go func() { _ = srv.Serve(ln) }() + t.Cleanup(func() { _ = srv.Close() }) + + conn, err := net.Dial("tcp", ln.Addr().String()) + require.NoError(t, err) + t.Cleanup(func() { _ = conn.Close() }) + + _, err = conn.Write([]byte("GET / HTTP/1.1\r\nHost: test\r\n\r\n")) + require.NoError(t, err) + + // Consume the full response so the connection becomes idle (keepalive). + br := bufio.NewReader(conn) + resp, err := http.ReadResponse(br, nil) + require.NoError(t, err) + _, err = io.ReadAll(resp.Body) + require.NoError(t, err) + require.NoError(t, resp.Body.Close()) + idleStart := time.Now() + + // Hold the connection open without another request. The server must close + // it within IdleTimeout + slack; the client observes the close as io.EOF. + require.NoError(t, conn.SetReadDeadline(time.Now().Add(testIdleTimeout+slack))) + _, err = br.Read(make([]byte, 1)) + elapsed := time.Since(idleStart) + + assert.ErrorIs(t, err, io.EOF, + "server should close the idle keepalive connection (got %v after %s)", err, elapsed) + assert.Less(t, elapsed, testIdleTimeout+slack, + "idle connection was not reclaimed within IdleTimeout+slack") +} - err := setUnixSocketPermission(path, tt.permission) - if tt.wantErr { - require.Error(t, err) - return +// TestSetupServerTimeoutsFieldValues pins the deliberate field decisions with +// no timed wait: ReadTimeout stays 0 on purpose (request +// bodies here are small; slow-loris is covered by ReadHeaderTimeout; SSE +// overrides WriteTimeout per-response) while the other three bound the wedge. +func TestSetupServerTimeoutsFieldValues(t *testing.T) { + t.Parallel() + + var srv http.Server + setupServerTimeouts(&srv) + + assert.Equal(t, time.Duration(0), srv.ReadTimeout, + "ReadTimeout must stay 0 (deliberate: small bodies, gRPC log upload is separate)") + assert.Equal(t, 10*time.Second, srv.ReadHeaderTimeout) + assert.Equal(t, 60*time.Second, srv.WriteTimeout) + assert.Equal(t, 120*time.Second, srv.IdleTimeout) +} + +// wantServerConstructions is how many http.Server literals this package builds: +// run() builds four across its two branches. Adding or removing a server is a +// deliberate act, so updating this number is part of it. +const wantServerConstructions = 4 + +// TestEveryHTTPServerGetsTimeouts pins the WIRING: every http.Server this +// package constructs must be handed to setupServerTimeouts. A helper that sets +// the right fields is worthless at a construction site that forgets to call it, +// and dropping one call is the realistic regression — run() builds four servers +// across two branches, and this file is a fork of upstream, so run() is rebased +// against upstream changes and a call is easy to lose in a conflict resolution. +// +// The run() function cannot be invoked from a test: it binds listeners, opens a +// store, and blocks. So the invariant is checked where it actually lives — in +// the source. +// The package is parsed and every composite literal of type http.Server is +// matched against the setupServerTimeouts arguments found in the same function. +// A new server added without the call fails here, naming the line. +// +// This is a stopgap for a structural gap: with each construction site extracted +// into a named constructor that ends in setupServerTimeouts, the same invariant +// would be a plain table test over those constructors' return values, and this +// AST walk could go away. +func TestEveryHTTPServerGetsTimeouts(t *testing.T) { + t.Parallel() + + fset := token.NewFileSet() + pkgs, err := parser.ParseDir(fset, ".", func(fi os.FileInfo) bool { //nolint:staticcheck // ParseDir is adequate for this build-tag-agnostic AST-only walk over the package's non-test sources + return !strings.HasSuffix(fi.Name(), "_test.go") + }, 0) + require.NoError(t, err) + require.NotEmpty(t, pkgs, "no non-test sources parsed; the walk below would vacuously pass") + + // Per enclosing function: where http.Server literals are built, and which + // identifiers were passed to setupServerTimeouts. + type siteSet struct { + built map[string]token.Position // variable name -> literal position + covered map[string]struct{} // names passed to setupServerTimeouts + } + + var totalBuilt int + for _, pkg := range pkgs { + for _, file := range pkg.Files { + for _, decl := range file.Decls { + fn, ok := decl.(*ast.FuncDecl) + if !ok || fn.Body == nil { + continue + } + // setupServerTimeouts itself takes the server as a parameter; + // it is the definition of the contract, not a call site. + if fn.Name.Name == "setupServerTimeouts" { + continue + } + + sites := siteSet{ + built: map[string]token.Position{}, + covered: map[string]struct{}{}, + } + ast.Inspect(fn.Body, func(n ast.Node) bool { + switch node := n.(type) { + case *ast.AssignStmt: + for i, rhs := range node.Rhs { + if !isHTTPServerLiteral(rhs) || i >= len(node.Lhs) { + continue + } + name, ok := node.Lhs[i].(*ast.Ident) + if !ok { + continue + } + sites.built[name.Name] = fset.Position(rhs.Pos()) + } + case *ast.ValueSpec: + // `var x = &http.Server{…}`. ast.Inspect descends into + // DeclStmt, so this needs no separate walk — but without + // this case the declaration form is invisible and a + // server built that way is silently unchecked. + for i, val := range node.Values { + if !isHTTPServerLiteral(val) || i >= len(node.Names) { + continue + } + sites.built[node.Names[i].Name] = fset.Position(val.Pos()) + } + case *ast.CallExpr: + id, ok := node.Fun.(*ast.Ident) + if !ok || id.Name != "setupServerTimeouts" || len(node.Args) != 1 { + return true + } + if arg, ok := node.Args[0].(*ast.Ident); ok { + sites.covered[arg.Name] = struct{}{} + } + } + return true + }) + + for name, pos := range sites.built { + totalBuilt++ + _, ok := sites.covered[name] + assert.True(t, ok, + "%s: http.Server %q is constructed in %s() without a "+ + "setupServerTimeouts call — its connections are unbounded", + pos, name, fn.Name.Name) + } } + } + } - require.NoError(t, err) - info, err := os.Stat(path) - require.NoError(t, err) - assert.Equal(t, tt.want, info.Mode().Perm()) - }) + // Guard the walk itself. NotZero is too weak: it still passes if the walk + // silently stops seeing three of the four sites — which is exactly what a + // construction written in a form the walk does not match looks like. Pin + // the count, so losing a site to an unmatched form fails here instead of + // passing over less code than it did yesterday. + assert.Equal(t, wantServerConstructions, totalBuilt, + "expected %d http.Server constructions in this package, walked %d — "+ + "either a server was added or removed (update this count), or one is "+ + "written in a form the walk does not match and is now unchecked", + wantServerConstructions, totalBuilt) +} + +// isHTTPServerLiteral reports whether expr builds an http.Server, as either +// &http.Server{…} or http.Server{…}. +func isHTTPServerLiteral(expr ast.Expr) bool { + if unary, ok := expr.(*ast.UnaryExpr); ok && unary.Op == token.AND { + expr = unary.X + } + lit, ok := expr.(*ast.CompositeLit) + if !ok { + return false + } + sel, ok := lit.Type.(*ast.SelectorExpr) + if !ok || sel.Sel.Name != "Server" { + return false } + pkg, ok := sel.X.(*ast.Ident) + return ok && pkg.Name == "http" } diff --git a/server/api/login.go b/server/api/login.go index c0ddf574a70..e1992f0e6b6 100644 --- a/server/api/login.go +++ b/server/api/login.go @@ -15,6 +15,7 @@ package api import ( + "context" "encoding/base32" "errors" "fmt" @@ -147,7 +148,17 @@ func HandleAuth(c *gin.Context) { allowedOrgs := server.Config.Permissions.Orgs.With(forgeModel.Orgs) if allowedOrgs.IsConfigured { isMember := false + // This paging blocks the login redirect and can outrun the server-wide + // WriteTimeout on an org-filtered account with many teams. Re-arm per + // page like updateRepoPermissions below, and stop paging if the browser + // gives up waiting. + onPage := slowHandlerProgress(c, newResponseController(c.Writer)) for page := 1; page <= maxPage; page++ { + if perr := onPage(); perr != nil { + log.Debug().Err(perr).Msgf("auth: client has gone away while verifying team membership for %s", userFromForge.Login) + c.AbortWithStatus(statusClientClosedRequest) + return + } teams, terr := _forge.Teams(c, userFromForge, &model.ListOptions{ Page: page, PerPage: perPage, @@ -315,16 +326,34 @@ func HandleAuth(c *gin.Context) { } if !server.Config.Server.AsyncRepositoryUpdate || noStoredRepositories { - if err := updateRepoPermissions(c, user, _store, _forge, forgeID); err != nil { - if err != nil { - log.Error().Err(err).Msgf("cannot update repo permissions for user %s", user.Login) + // Synchronous: the redirect below is blocked on this sync, which pages + // against the forge and can outrun the server-wide WriteTimeout. Re-arm + // per page so a login against a large account still completes, and stop + // paging if the browser gives up waiting for the redirect. + onPage := slowHandlerProgress(c, newResponseController(c.Writer)) + if err := updateRepoPermissions(c, user, _store, _forge, forgeID, onPage); err != nil { + // onPage reports the browser gave up waiting for the redirect. There + // is nobody left to redirect, and a login abandoned mid-sync is not + // an internal error: log it at debug and let the access log record + // the disconnect rather than a 500 nobody caused. + if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) { + log.Debug().Err(err).Msgf("auth: client has gone away while syncing repo permissions for user %s", user.Login) + c.AbortWithStatus(statusClientClosedRequest) + return } + log.Error().Err(err).Msgf("cannot update repo permissions for user %s", user.Login) c.Redirect(http.StatusSeeOther, server.Config.Server.RootPath+"/login?error=internal_error") return } } else { go func() { - if err := updateRepoPermissions(c, user, _store, _forge, forgeID); err != nil { + // Detached from the response: this outlives the handler, which + // writes its redirect below and returns. Touching the response + // writer from here would race that return, so there is no deadline + // of ours to extend — hence the nil hook, and hence why the sync + // takes a progress callback rather than a ResponseController, which + // only its HTTP callers can supply. + if err := updateRepoPermissions(c, user, _store, _forge, forgeID, nil); err != nil { log.Error().Err(err).Msgf("could not update repo permissions for user %s in background", user.Login) } }() @@ -334,9 +363,27 @@ func HandleAuth(c *gin.Context) { c.Redirect(http.StatusSeeOther, server.Config.Server.RootPath+"/") } -func updateRepoPermissions(c *gin.Context, user *model.User, _store store.Store, _forge forge.Forge, forgeID int64) error { +// updateRepoPermissions syncs a user's forge repo permissions into the store. +// +// The onPage hook, when non-nil, runs before fetching each page from the forge +// and can abort the sync by returning an error. It is the seam an HTTP caller +// uses both to re-arm its rolling per-response write deadline as the sync makes +// progress and to stop paging once its client has gone away — necessary because +// rolling the deadline removes the wall-clock bound that would otherwise end the +// handler, and maxPage alone permits 10000 forge round-trips against a response +// nobody is reading. +// +// The background caller passes nil: it has no live response to arm and must +// keep running detached. A callback rather than a ResponseController keeps this +// function free of a dependency only some of its callers can satisfy. +func updateRepoPermissions(c *gin.Context, user *model.User, _store store.Store, _forge forge.Forge, forgeID int64, onPage func() error) error { start := time.Now() repos, err := utils.Paginate(func(page int) ([]*model.Repo, error) { + if onPage != nil { + if err := onPage(); err != nil { + return nil, err + } + } return _forge.Repos(c, user, &model.ListOptions{ Page: page, PerPage: perPage, diff --git a/server/api/login_test.go b/server/api/login_test.go index b101edec5cf..05b97922569 100644 --- a/server/api/login_test.go +++ b/server/api/login_test.go @@ -472,6 +472,84 @@ func TestHandleAuth(t *testing.T) { }) } +// The synchronous permission sync in HandleAuth gained a 499-vs-internal_error +// split: a browser that gives up mid-sync is logged as a client disconnect and +// aborted with 499, while a genuine forge error still redirects to the login +// page with error=internal_error. The slowHandlerProgress hook checks the +// context before each forge page, so a request whose context is already +// canceled makes the sync return context.Canceled — the same error the abort +// branch keys on — without depending on real disconnect timing. +func TestHandleAuthClassifiesClientCancelVsForgeError(t *testing.T) { + gin.SetMode(gin.TestMode) + + user := &model.User{ + ID: 1, + OrgID: 1, + ForgeID: 1, + ForgeRemoteID: "remote-id-1", + Login: "test", + Email: "test@example.com", + } + org := &model.Org{ID: 1, Name: user.Login} + server.Config.Server.SessionExpires = time.Hour + + // wireAuthedLogin sets up the mock chain for a login that reaches the + // synchronous permission sync: RepoList returns no stored repos, so the + // sync runs inline rather than in the background. + wireAuthedLogin := func(t *testing.T) (*forge_mocks.MockForge, *store_mocks.MockStore) { + t.Helper() + _manager := manager_mocks.NewMockManager(t) + _forge := forge_mocks.NewMockForge(t) + _store := store_mocks.NewMockStore(t) + server.Config.Services.Manager = _manager + server.Config.Permissions.Open = true + server.Config.Permissions.Orgs = permissions.NewOrgs(nil) + server.Config.Permissions.Admins = permissions.NewAdmins(nil) + + _manager.On("ForgeByID", int64(1)).Return(_forge, nil) + _store.On("ForgeGet", int64(1)).Return(&model.Forge{ID: 1}, nil) + _forge.On("Login", mock.Anything, mock.Anything).Return(user, "", nil) + _store.On("GetUserByRemoteID", user.ForgeID, user.ForgeRemoteID).Return(user, nil) + _store.On("OrgGet", org.ID).Return(org, nil) + _store.On("UpdateUser", mock.Anything).Return(nil) + // empty repo list => noStoredRepositories, so the sync runs synchronously + _store.On("RepoList", mock.Anything, mock.Anything, mock.Anything, mock.Anything).Return(nil, nil) + return _forge, _store + } + + t.Run("client disconnect during sync aborts with 499", func(t *testing.T) { + _, _store := wireAuthedLogin(t) + w := httptest.NewRecorder() + c, _ := gin.CreateTestContext(w) + c.Set("store", _store) + ctx, cancel := context.WithCancelCause(context.Background()) + cancel(nil) + c.Request = httptest.NewRequestWithContext(ctx, http.MethodGet, "https://example.com/authorize", nil) + + api.HandleAuth(c) + + // A canceled request trips slowHandlerProgress before the first forge + // page, so the sync returns context.Canceled and the handler aborts with + // nginx's 499 (client closed request) rather than redirecting the + // browser that already left. + assert.Equal(t, 499, c.Writer.Status()) + }) + + t.Run("forge error during sync redirects to internal_error", func(t *testing.T) { + _forge, _store := wireAuthedLogin(t) + _forge.On("Repos", mock.Anything, mock.Anything, mock.Anything).Return(nil, assert.AnError) + w := httptest.NewRecorder() + c, _ := gin.CreateTestContext(w) + c.Set("store", _store) + c.Request = httptest.NewRequest(http.MethodGet, "https://example.com/authorize", nil) + + api.HandleAuth(c) + + assert.Equal(t, http.StatusSeeOther, c.Writer.Status()) + assert.Equal(t, "/login?error=internal_error", c.Writer.Header().Get("Location")) + }) +} + // TestHandleAuthAllowedOrgs walks through the combinations of WOODPECKER_ORGS // and the orgs of the forge a user logs in with. The global list applies to // every forge, the orgs of a forge are allowed in addition to it. Org names are diff --git a/server/api/queue.go b/server/api/queue.go index e9ea4d53d93..a006fea7194 100644 --- a/server/api/queue.go +++ b/server/api/queue.go @@ -17,6 +17,7 @@ package api import ( "fmt" "net/http" + "time" "github.com/gin-gonic/gin" "github.com/rs/zerolog/log" @@ -26,6 +27,31 @@ import ( "go.woodpecker-ci.org/woodpecker/v3/server/store" ) +// slowHandlerWriteExtension is the per-response write-deadline extension taken +// by the handlers in this package that can legitimately run longer than the +// server-wide WriteTimeout — the queue long-poll below, the all-repo repair, +// and the login-time permission sync, each of which is bounded by external +// work (a scheduler draining, a forge answering) rather than by its own +// compute. See writedeadline.go for why that budget covers handler runtime and +// not just blocked writes. +// +// It is sized at roughly the server-wide WriteTimeout on purpose: a handler +// that keeps making progress re-arms before the previous extension lapses and +// so never trips, while one that wedges between two re-arm points still dies +// on a budget no larger than the server-wide one. Making it much shorter would +// trip handlers whose next unit of work is a slow-but-healthy forge call; +// making it much longer would let a wedged handler hold a connection well past +// what the server-wide timeout promises. +const slowHandlerWriteExtension = 60 * time.Second + +// queuePollInterval paces the queue long-poll below. The loop used to spin hot +// on the scheduler, burning a core for the whole wait; it now sleeps between +// probes, which also gives it a sane cadence at which to re-arm its write +// deadline instead of issuing that syscall millions of times a second. Short +// enough that the observed latency of the 204 is dominated by the queue +// draining, not by the poll. +const queuePollInterval = 100 * time.Millisecond + // GetQueueInfo // // @Summary Get pipeline queue information @@ -118,11 +144,34 @@ func ResumeQueue(c *gin.Context) { // @Tags Pipeline queues // @Param Authorization header string true "Insert your personal access token" default(Bearer ) func BlockTilQueueHasRunningItem(c *gin.Context) { + // This long-poll returns only once the queue has drained, which can take far + // longer than the server-wide WriteTimeout. That budget covers total handler + // runtime, so without a rolling per-response deadline the connection is torn + // down mid-wait and the documented 204 never reaches the client. + // + // Rolling the deadline removes the only bound this handler had, so it must + // supply its own: the loop writes nothing, so no write can ever trip the + // deadline and end it. The request context is the replacement — a client + // hang-up ends the wait immediately. Without it, every abandoned request + // would pin a goroutine and a 10 Hz poll on the scheduler lock until the + // process restarted. + progress := slowHandlerProgress(c, newResponseController(c.Writer)) for { + // A canceled request ends the wait: there is nobody to send the 204 to. + if err := progress(); err != nil { + return + } + info := server.Config.Services.Scheduler.Info(c) if info.Stats.Running == 0 { break } + + select { + case <-c.Request.Context().Done(): + return + case <-time.After(queuePollInterval): + } } c.Status(http.StatusNoContent) } diff --git a/server/api/queue_test.go b/server/api/queue_test.go index 729cbdb94a3..e8ee6326655 100644 --- a/server/api/queue_test.go +++ b/server/api/queue_test.go @@ -17,9 +17,14 @@ package api import ( + "context" + "net" "net/http" + "sync" "testing" + "time" + "github.com/gin-gonic/gin" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/mock" "github.com/stretchr/testify/require" @@ -191,3 +196,91 @@ func TestPauseResumeQueue(t *testing.T) { q.AssertCalled(t, "Resume") }) } + +// BlockTilQueueHasRunningItem waits for a condition that may never arrive, and +// arming the rolling write deadline removed the bound that used to end it: the +// loop writes nothing, so no write can trip a deadline, and WriteTimeout no +// longer applies. Client disconnect is the only remaining bound, so this pins +// it directly — without it the handler goroutine runs until the process does. +func TestBlockTilQueueHasRunningItemStopsWhenTheClientDisconnects(t *testing.T) { + gin.SetMode(gin.TestMode) + + s := newTestStore(t) + q := installScheduler(t, s) + // Never drains: Running stays at 1 for the life of the test, so the only + // way out of the loop is the client going away. polled closes on the first + // Info call, gating the hangup on the handler actually reaching its loop — + // reading the mock's call log from the test goroutine would race the + // handler's writes to it. + info := queue.InfoT{} + info.Stats.Running = 1 + polled := make(chan struct{}) + var polledOnce sync.Once + q.On("Info", mock.Anything).Return(info). + Run(func(mock.Arguments) { polledOnce.Do(func() { close(polled) }) }) + + handlerDone := make(chan struct{}) + router := gin.New() + router.GET("/queue/resume", BlockTilQueueHasRunningItem) + + // Signal after gin has fully finished with the request, not from inside the + // handler: gin writes the response header after the handler returns, so a + // defer inside it would release the test while gin still owns the context + // the mock's cleanup formats. + handler := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + defer close(handlerDone) + router.ServeHTTP(w, r) + }) + + // A plain http.Server rather than httptest.Server: httptest's Close waits + // for outstanding handlers, so when this test fails — the handler leaked — + // cleanup would block forever and take the whole package's suite down with + // it. Closing the listener alone lets a failure stay a failed test. + ln, err := net.Listen("tcp", "127.0.0.1:0") + require.NoError(t, err) + // WriteTimeout unset: cancellation is the only bound under test. + srv := &http.Server{Handler: handler, ReadHeaderTimeout: 5 * time.Second} + go func() { _ = srv.Serve(ln) }() + t.Cleanup(func() { + // Close (not Shutdown) so a leaked handler cannot block cleanup, then + // give the handler a bounded moment to return. Without this wait a + // handler still in its poll loop outlives the test and calls the mock + // queue while a later test owns it. + _ = srv.Close() + select { + case <-handlerDone: + case <-time.After(5 * time.Second): + } + }) + + reqCtx, cancelReq := context.WithCancelCause(t.Context()) + url := "http://" + ln.Addr().String() + "/queue/resume" + req, err := http.NewRequestWithContext(reqCtx, http.MethodGet, url, nil) + require.NoError(t, err) + + reqDone := make(chan struct{}) + go func() { + defer close(reqDone) + resp, err := http.DefaultClient.Do(req) + if err == nil { + _ = resp.Body.Close() + } + }() + + // Let the handler reach its loop before hanging up, so the test exercises + // cancellation mid-wait rather than a request that never started. + select { + case <-polled: + case <-time.After(10 * time.Second): + t.Fatal("handler never polled the queue") + } + + cancelReq(nil) + <-reqDone + + select { + case <-handlerDone: + case <-time.After(10 * time.Second): + t.Fatal("handler outlived its client: the long-poll has no cancellation bound and leaks a goroutine per hangup") + } +} diff --git a/server/api/repo.go b/server/api/repo.go index ddc4536a259..7aa171ba58e 100644 --- a/server/api/repo.go +++ b/server/api/repo.go @@ -710,8 +710,28 @@ func RepairAllRepos(c *gin.Context) { return } + // Each repairRepo below is at least one forge round-trip, so a large install + // easily outruns the server-wide WriteTimeout, which bounds total handler + // runtime. Re-arm once per repo: steady progress keeps the response alive, + // while a repair wedged on an unresponsive forge still dies on the original + // budget. + progress := slowHandlerProgress(c, newResponseController(c.Writer)) + failedRepos := make([]int64, 0) for _, r := range repos { + // Rolling the deadline means no write can end this loop early, and each + // repair is several forge round-trips. Stop when the client goes away + // rather than working through the whole install for a response nobody + // will read. + if err := progress(); err != nil { + // A partial repair is not a success. Without an explicit status gin + // finalizes a bare 200 with an empty body and the access log records + // the abandoned repair as a completed one. + log.Debug().Err(err).Msg("repair all repos: client has gone away, stopping") + c.AbortWithStatus(statusClientClosedRequest) + return + } + // updatePermissions is false as RepoListAll does not load permissions updatePermissions := false err := repairRepo(c, r, updatePermissions) diff --git a/server/api/stream.go b/server/api/stream.go index 1451ce4c5cd..e6a4b8e701c 100644 --- a/server/api/stream.go +++ b/server/api/stream.go @@ -20,6 +20,7 @@ import ( "encoding/json" "errors" "io" + "math/rand/v2" "net/http" "strconv" "sync" @@ -40,11 +41,41 @@ const ( // How many batches of logs to keep for each client before starting to // drop them if the client is not consuming them faster than they arrive. maxQueuedBatchesPerClient int = 30 - - // Is the time till we send a ping to keep the connection alive. - idlePingTime = time.Second * 30 ) +// idlePingTime is the interval between keep-alive pings on an SSE stream. It is +// a var, not a const, so tests can shorten it to exercise the rolling +// SetWriteDeadline path without real-time waits. Tests that mutate it MUST NOT +// run with t.Parallel(): it is a shared package-level global read by the SSE +// handlers, so a parallel mutator races the handler's read under `go test -race`. +var idlePingTime = time.Second * 30 + +// streamMaxDuration is the absolute ceiling on a single SSE response, measured +// from the start of the request. It backstops the rolling write deadline for a +// stream whose peer stops reading without closing and which carries only +// keep-alive pings — see armSSEWriteDeadline. Reaching it ends the response; +// the client's EventSource reconnects. +// +// A var for the same test reason as idlePingTime, and under the same rule. +var streamMaxDuration = time.Hour + +// streamCeilingJitterDivisor bounds the ceiling jitter to streamMaxDuration +// divided by this factor (i.e. up to a tenth of the ceiling). See streamCeiling. +const streamCeilingJitterDivisor = 10 + +// streamCeiling returns the absolute per-connection deadline for one SSE +// response: now + streamMaxDuration, less a bounded random jitter in +// [0, streamMaxDuration/10]. The jitter de-synchronizes a cohort of clients +// that connected together (e.g. every dashboard tab reconnecting at once after +// a server restart), so their ceilings expire spread out rather than re-forming +// a thundering herd every streamMaxDuration. This is load spreading rather than +// a security boundary, so math/rand/v2 is sufficient. Computed once per handler +// invocation. +func streamCeiling() time.Time { + jitter := rand.N(streamMaxDuration / streamCeilingJitterDivisor) //nolint:gosec // load spreading, not security + return time.Now().Add(streamMaxDuration - jitter) +} + // EventStreamSSE // // @Summary Stream events like pipeline updates @@ -61,15 +92,27 @@ func EventStreamSSE(c *gin.Context) { rw := c.Writer - flusher, ok := rw.(http.Flusher) - if !ok { - c.String(http.StatusInternalServerError, "Streaming not supported") + rc := newResponseController(rw) + // One warn per response if this writer cannot take a deadline; the arms + // below run per write. + var deadlineWarnOnce sync.Once + streamUntil := streamCeiling() + + // ping the client. Arm before the write, as every write site does — the + // first flush is also what reports a writer that cannot stream at all, + // since rc.Flush surfaces ErrNotSupported where a Flusher type-assertion + // on the gin writer would not (gin always implements Flush; the wrapped + // writer is the one that matters). + armSSEWriteDeadline(rc, streamUntil, &deadlineWarnOnce) + if _, err := io.WriteString(rw, ": ping\n\n"); err != nil { + return + } + if err := rc.Flush(); err != nil { + if errors.Is(err, http.ErrNotSupported) { + c.String(http.StatusInternalServerError, "Streaming not supported") + } return } - - // ping the client - logWriteStringErr(io.WriteString(rw, ": ping\n\n")) - flusher.Flush() log.Debug().Msg("user feed: connection opened") @@ -96,8 +139,11 @@ func EventStreamSSE(c *gin.Context) { log.Debug().Msg("user feed: connection closed") }() + // Captured once: this goroutine outlives the handler, and reading the global + // from inside it would read it after a test's cleanup has reset it. + scheduler := server.Config.Services.Scheduler go func() { - err := server.Config.Services.Scheduler.Subscribe(ctx, subTopics, + err := scheduler.Subscribe(ctx, subTopics, func(m pubsub.Message) { select { case <-ctx.Done(): @@ -107,21 +153,44 @@ func EventStreamSSE(c *gin.Context) { cancel(err) }() + streamExpiry := time.NewTimer(time.Until(streamUntil)) + defer streamExpiry.Stop() + for { select { case <-requestCtx.Done(): return case <-ctx.Done(): return + case <-streamExpiry.C: + // The absolute ceiling: a stream that reached it has run its full + // allowance without the peer closing. Ending the response releases + // the handler and subscription; a live client reconnects. + log.Debug().Msg("user feed: stream reached its maximum duration") + return case <-time.After(idlePingTime): - logWriteStringErr(io.WriteString(rw, ": ping\n\n")) - flusher.Flush() + armSSEWriteDeadline(rc, streamUntil, &deadlineWarnOnce) + if _, err := io.WriteString(rw, ": ping\n\n"); err != nil { + return + } + if err := rc.Flush(); err != nil { + return + } case buf, ok := <-eventChan: if ok { - logWriteStringErr(io.WriteString(rw, "data: ")) - logWriteStringErr(rw.Write(buf)) - logWriteStringErr(io.WriteString(rw, "\n\n")) - flusher.Flush() + armSSEWriteDeadline(rc, streamUntil, &deadlineWarnOnce) + if _, err := io.WriteString(rw, "data: "); err != nil { + return + } + if _, err := rw.Write(buf); err != nil { + return + } + if _, err := io.WriteString(rw, "\n\n"); err != nil { + return + } + if err := rc.Flush(); err != nil { + return + } } } } @@ -145,14 +214,26 @@ func LogStreamSSE(c *gin.Context) { rw := c.Writer - flusher, ok := rw.(http.Flusher) - if !ok { - c.String(http.StatusInternalServerError, "Streaming not supported") + rc := newResponseController(rw) + // One warn per response if this writer cannot take a deadline; the arms + // below run per write. + var deadlineWarnOnce sync.Once + streamUntil := streamCeiling() + + // Arm before the write, as every write site does. The first flush doubles + // as the streaming-capability check: rc.Flush surfaces ErrNotSupported from + // the writer that actually carries the flush, which a Flusher assertion on + // the gin wrapper would not. + armSSEWriteDeadline(rc, streamUntil, &deadlineWarnOnce) + if _, err := io.WriteString(rw, ": ping\n\n"); err != nil { + return + } + if err := rc.Flush(); err != nil { + if errors.Is(err, http.ErrNotSupported) { + c.String(http.StatusInternalServerError, "Streaming not supported") + } return } - - logWriteStringErr(io.WriteString(rw, ": ping\n\n")) - flusher.Flush() _store := store.FromContext(c) repo := session.Repo(c) @@ -202,7 +283,10 @@ func LogStreamSSE(c *gin.Context) { log.Debug().Msg("log stream: connection closed") }() - err = server.Config.Services.Logs.Open(ctx, step.ID) + // Captured once: the tail goroutine below outlives this handler, and reading + // the global from there would read it after a test's cleanup has reset it. + logService := server.Config.Services.Logs + err = logService.Open(ctx, step.ID) if err != nil { log.Error().Err(err).Msg("log stream: open failed") logWriteStringErr(io.WriteString(rw, "event: error\ndata: can't open stream\n\n")) @@ -231,7 +315,7 @@ func LogStreamSSE(c *gin.Context) { } }() - err := server.Config.Services.Logs.Tail(ctx, step.ID, batches) + err := logService.Tail(ctx, step.ID, batches) if err != nil { log.Error().Err(err).Msg("tail of logs failed") } @@ -249,30 +333,77 @@ func LogStreamSSE(c *gin.Context) { log.Debug().Msgf("log stream: reconnect: last-event-id: %d", last) } + streamExpiry := time.NewTimer(time.Until(streamUntil)) + defer streamExpiry.Stop() + for { select { case <-ctx.Done(): // Monitor if the "tail" context is canceled. + // Return UNCONDITIONALLY: the tail context is done, so this stream + // is over regardless of cause. Only the eof marker is conditional — + // an ordinary end-of-logs cancellation (context.Canceled) tells the + // client to stop; any other cause (e.g. a Tail error like ErrNotFound + // on a race) has no marker to send. Falling through without returning + // would re-select an already-closed ctx.Done() and spin at 100% CPU + // until the ceiling. if err := context.Cause(ctx); errors.Is(err, context.Canceled) { log.Debug().Msg("log stream: eof") + // Arm before this write like every other write site. The last + // arm can be arbitrarily stale by now: any select arm firing + // restarts the ping countdown, and the replay path re-arms only + // inside `id > last`, so a reconnect that replays for longer + // than the deadline leaves it already expired. Without a fresh + // arm a healthy, still-reading client silently never receives + // the eof marker and its EventSource reconnects forever. + armSSEWriteDeadline(rc, streamUntil, &deadlineWarnOnce) logWriteStringErr(io.WriteString(rw, "event: eof\ndata: eof\n\n")) - flusher.Flush() - return + // Best-effort: the stream is ending on this return either way, + // so a failed final flush only means the peer never saw the + // eof marker. Logged rather than discarded silently. + if err := rc.Flush(); err != nil { + log.Debug().Err(err).Msg("log stream: flushing eof marker") + } } + return case <-requestCtx.Done(): // Monitor the request context for cancellation when the client has gone away. log.Debug().Msg("log stream: closed, client has gone away") return + case <-streamExpiry.C: + // The absolute ceiling — see armSSEWriteDeadline. A log stream that + // reaches it ends without an eof marker, because the logs have not + // ended; the client reconnects with Last-Event-ID and resumes. + log.Debug().Msg("log stream: reached its maximum duration") + return case <-time.After(idlePingTime): - logWriteStringErr(io.WriteString(rw, ": ping\n\n")) - flusher.Flush() + armSSEWriteDeadline(rc, streamUntil, &deadlineWarnOnce) + if _, err := io.WriteString(rw, ": ping\n\n"); err != nil { + return + } + if err := rc.Flush(); err != nil { + return + } case buf, ok := <-logChan: if ok { if id > last { - logWriteStringErr(io.WriteString(rw, "id: "+strconv.Itoa(id))) - logWriteStringErr(io.WriteString(rw, "\n")) - logWriteStringErr(io.WriteString(rw, "data: ")) - logWriteStringErr(rw.Write(buf)) - logWriteStringErr(io.WriteString(rw, "\n\n")) - flusher.Flush() + armSSEWriteDeadline(rc, streamUntil, &deadlineWarnOnce) + if _, err := io.WriteString(rw, "id: "+strconv.Itoa(id)); err != nil { + return + } + if _, err := io.WriteString(rw, "\n"); err != nil { + return + } + if _, err := io.WriteString(rw, "data: "); err != nil { + return + } + if _, err := rw.Write(buf); err != nil { + return + } + if _, err := io.WriteString(rw, "\n\n"); err != nil { + return + } + if err := rc.Flush(); err != nil { + return + } } id++ } @@ -285,3 +416,45 @@ func logWriteStringErr(_ int, err error) { log.Error().Err(err).Caller(1).Msg("fail to write string") } } + +// armSSEWriteDeadline arms the rolling per-response write deadline on an SSE +// stream, refreshed before each write so a healthy stream outlives the +// server-wide WriteTimeout while a peer that stops reading still trips the +// deadline and unblocks the handler. +// +// A ping-only stream is reclaimed by the ceiling TIMER, not by this arm: each +// handler builds a streamExpiry timer for streamMaxDuration and returns on its +// arm (EventStreamSSE, LogStreamSSE). That timer is what bounds the idle +// /api/stream/events case every open browser tab holds — the rolling deadline +// alone would not, because it only fires when a write actually blocks, and a +// 9-byte keep-alive ping never fills a socket send buffer (roughly 277k pings, +// some 96 days at the default interval, before one would). +// +// The clamp below is narrower: it keeps an arm from landing in the past as the +// response approaches streamMaxDuration, so a write racing the expiry timer +// ends the stream on the timer's clean return rather than on a write error. +// +// A healthy stream is unaffected until the ceiling: it re-arms every +// idlePingTime against a 2*idlePingTime deadline, so it keeps a full +// idlePingTime of slack. On reaching the ceiling the handler returns and the +// client's EventSource reconnects, which is the normal SSE lifecycle. +// +// The controller passed here must be one built by newResponseController: the +// ping's only socket write happens inside the flush whose error gin would +// otherwise swallow, so a ping-only stream is the case that most needs the +// unwrap. +func armSSEWriteDeadline(rc *http.ResponseController, until time.Time, warnOnce *sync.Once) { + // A healthy stream re-arms every idlePingTime against a deadline two ping + // intervals out, so it always keeps a full idlePingTime of slack. + const pingSlackIntervals = 2 + d := pingSlackIntervals * idlePingTime + // Clamp to the ceiling, but never to a deadline in the past. A write can + // race the expiry timer and find no time remaining; arming that verbatim + // would fail the write and end the stream on an error rather than on the + // clean return the expiry arm makes. The floor lets the in-flight write + // finish — the timer ends the stream on the next loop either way. + if remaining := time.Until(until); remaining < d { + d = max(remaining, idlePingTime) + } + extendWriteDeadline(rc, d, warnOnce) +} diff --git a/server/api/stream_test.go b/server/api/stream_test.go index d4a7032c975..534fee972d3 100644 --- a/server/api/stream_test.go +++ b/server/api/stream_test.go @@ -15,16 +15,22 @@ package api import ( + "bufio" + "bytes" "context" "fmt" + "io" "net/http" "net/http/httptest" + "os" + "strings" "sync" "testing" "time" "github.com/gin-gonic/gin" "github.com/stretchr/testify/mock" + "github.com/stretchr/testify/require" "go.woodpecker-ci.org/woodpecker/v3/server" "go.woodpecker-ci.org/woodpecker/v3/server/logging" @@ -158,3 +164,845 @@ func TestLogStreamSSEConcurrentDisconnect(t *testing.T) { }) } } + +// newSSEEventServer wires the real EventStreamSSE handler behind a gin router, +// backed by an in-memory pubsub broker, and returns a running httptest server. +// The handlerDone channel is closed when the handler returns; start selects the +// transport (plain HTTP/1 or HTTP/2+TLS) so the same assertion runs over both. +func newSSEEventServer( + t *testing.T, + writeTimeout time.Duration, + start func(*httptest.Server), +) (ts *httptest.Server, broker pubsub.PubSub, handlerDone <-chan struct{}) { + t.Helper() + + broker = memory.New() + server.Config.Services.Scheduler = scheduler.NewScheduler(t.Context(), nil, nil, broker) + t.Cleanup(func() { server.Config.Services.Scheduler = nil }) + + done := make(chan struct{}) + router := gin.New() + router.GET("/stream/events", func(c *gin.Context) { + defer close(done) + EventStreamSSE(c) + }) + + ts = httptest.NewUnstartedServer(router) + ts.Config.WriteTimeout = writeTimeout + start(ts) + t.Cleanup(ts.Close) + + return ts, broker, done +} + +// TestEventStreamSSESurvivesWriteTimeout pins the rolling-deadline override: a +// stream keeps delivering events well past the server-wide WriteTimeout, +// because the handler re-arms a rolling per-response write deadline +// (2*idlePingTime) before every write. It runs over both HTTP/1 and HTTP/2+TLS +// so the h2 responseWriter's SetWriteDeadline path is exercised too. Without +// the rolling deadline the server would tear the stream down at WriteTimeout +// and no event would arrive after it — the assertion below would fail. +func TestEventStreamSSESurvivesWriteTimeout(t *testing.T) { + gin.SetMode(gin.TestMode) + + // The override proof needs the rolling deadline (2*idlePingTime) to outlast + // the server-wide WriteTimeout on the FIRST arm, so the stream survives past + // it even if the handler never gets to re-arm — otherwise survival would hinge + // on re-arming within the window every time, which a scheduler stall under + // load can miss (flaky). Keep 2*idlePingTime (200ms) > testWriteTimeout (80ms) + // with margin. Mutates the shared idlePingTime global, so no t.Parallel(). + origPing := idlePingTime + idlePingTime = 100 * time.Millisecond + t.Cleanup(func() { idlePingTime = origPing }) + + const testWriteTimeout = 80 * time.Millisecond + + transports := []struct { + name string + start func(*httptest.Server) + }{ + {"http1", func(ts *httptest.Server) { ts.Start() }}, + {"http2", func(ts *httptest.Server) { ts.EnableHTTP2 = true; ts.StartTLS() }}, + } + + for _, tr := range transports { + t.Run(tr.name, func(t *testing.T) { + ts, broker, handlerDone := newSSEEventServer(t, testWriteTimeout, tr.start) + + ctx, cancel := context.WithTimeout(t.Context(), 5*time.Second) + t.Cleanup(cancel) + + req, err := http.NewRequestWithContext(ctx, http.MethodGet, ts.URL+"/stream/events", nil) + require.NoError(t, err) + resp, err := ts.Client().Do(req) + require.NoError(t, err) + t.Cleanup(func() { _ = resp.Body.Close() }) + require.Equal(t, http.StatusOK, resp.StatusCode) + + // Emit data events continuously, past the WriteTimeout window. + pubStop := make(chan struct{}) + go func() { + ticker := time.NewTicker(30 * time.Millisecond) + defer ticker.Stop() + topic := map[string]struct{}{pubsub.PublicTopic: {}} + for { + select { + case <-pubStop: + return + case <-ticker.C: + _ = broker.Publish(context.Background(), topic, + pubsub.Message{Data: []byte(`{"pipeline":1}`)}) + } + } + }() + t.Cleanup(func() { close(pubStop) }) + + // Read the stream and confirm at least one line arrives strictly + // after WriteTimeout has elapsed — proof the rolling deadline + // overrode the server-wide one. + start := time.Now() + br := bufio.NewReader(resp.Body) + readDeadline := time.Now().Add(2*testWriteTimeout + time.Second) + var gotLate bool + for time.Now().Before(readDeadline) { + line, err := br.ReadString('\n') + if err != nil { + break + } + if strings.TrimSpace(line) == "" { + continue + } + if time.Since(start) > testWriteTimeout { + gotLate = true + break + } + } + require.True(t, gotLate, + "SSE stream stopped delivering before WriteTimeout elapsed; the rolling write deadline did not override the server-wide WriteTimeout") + + // Deterministically reap the handler goroutine before this subtest + // returns. Canceling the request trips the handler's + // requestCtx.Done() arm; waiting on handlerDone establishes a + // happens-before with the parent's idlePingTime restore, so the + // live handler can no longer race the restore's write of the + // global. This is the fix for that race — not a reason to rerun + // the test until it passes. + cancel() + select { + case <-handlerDone: + case <-time.After(2*testWriteTimeout + time.Second): + t.Fatal("SSE handler did not return after request cancel; goroutine leaked") + } + }) + } +} + +// TestEventStreamSSEWedgeKill pins the return-on-write-error path: when a +// client stops reading, the send buffer fills, and the blocked write is +// reclaimed by the rolling deadline (2*idlePingTime) — the handler observes the +// write error and returns instead of spinning forever. WriteTimeout is left 0 +// so the ONLY thing that can unblock the wedged write is the per-response +// deadline the fix arms; without it the handler would block indefinitely and +// this test would time out. +func TestEventStreamSSEWedgeKill(t *testing.T) { + gin.SetMode(gin.TestMode) + + origPing := idlePingTime + idlePingTime = 40 * time.Millisecond + t.Cleanup(func() { idlePingTime = origPing }) + + // WriteTimeout 0: no server-wide deadline. Only armSSEWriteDeadline can + // reclaim the wedged write. + ts, broker, handlerDone := newSSEEventServer(t, 0, func(ts *httptest.Server) { ts.Start() }) + + ctx, cancel := context.WithCancelCause(t.Context()) + t.Cleanup(func() { cancel(nil) }) + + req, err := http.NewRequestWithContext(ctx, http.MethodGet, ts.URL+"/stream/events", nil) + require.NoError(t, err) + resp, err := ts.Client().Do(req) + require.NoError(t, err) + t.Cleanup(func() { _ = resp.Body.Close() }) + require.Equal(t, http.StatusOK, resp.StatusCode) + + // The client deliberately never reads resp.Body. Flood large payloads so + // the kernel send buffer fills and the next handler write blocks. + pubStop := make(chan struct{}) + go func() { + ticker := time.NewTicker(5 * time.Millisecond) + defer ticker.Stop() + topic := map[string]struct{}{pubsub.PublicTopic: {}} + big := bytes.Repeat([]byte("x"), 256<<10) + for { + select { + case <-pubStop: + return + case <-ticker.C: + _ = broker.Publish(context.Background(), topic, pubsub.Message{Data: big}) + } + } + }() + t.Cleanup(func() { close(pubStop) }) + + // The wedged write must be reclaimed by the rolling deadline and the + // handler must return within 2*idlePingTime + slack. + const slack = 3 * time.Second + select { + case <-handlerDone: + // handler returned — the write error was surfaced and acted on. + case <-time.After(2*idlePingTime + slack): + t.Fatalf("SSE handler did not return within 2*idlePingTime+slack (%s); "+ + "the wedged write was never reclaimed / the write error was not acted on", 2*idlePingTime+slack) + } +} + +// newSSELogServer wires the real LogStreamSSE handler behind a gin router, +// backed by an in-memory logging service, and returns a running httptest +// server. Unlike setupLogStreamContext (which uses an httptest.NewRecorder that +// never blocks), this drives the handler over a real net.Conn so a real write +// can actually block and the rolling write deadline can fire. The route wrapper +// populates the context the way middleware would before calling LogStreamSSE: +// it sets the mock store (GetPipelineNumber -> Pipeline{ID:pipelineID}, +// StepLoad -> running Step{ID:stepID}) and repo, matching the route params +// repo_id=1/pipeline=1/step_id=42. The handlerDone channel is closed when the +// handler returns; start selects the transport (plain HTTP/1 or HTTP/2+TLS). +func newSSELogServer( + t *testing.T, + writeTimeout time.Duration, + start func(*httptest.Server), +) (ts *httptest.Server, logService logging.Log, handlerDone <-chan struct{}) { + t.Helper() + + const stepID int64 = 42 + const pipelineID int64 = 10 + + logService = logging.New() + server.Config.Services.Logs = logService + t.Cleanup(func() { server.Config.Services.Logs = nil }) + + mockStore := store_mocks.NewMockStore(t) + mockStore.On("GetPipelineNumber", mock.Anything, mock.Anything). + Return(&model.Pipeline{ID: pipelineID}, nil) + mockStore.On("StepLoad", mock.Anything, mock.Anything). + Return(&model.Step{ + ID: stepID, + PipelineID: pipelineID, + State: model.StatusRunning, + }, nil) + + done := make(chan struct{}) + router := gin.New() + router.GET("/stream/logs/:repo_id/:pipeline/:step_id", func(c *gin.Context) { + defer close(done) + c.Set("store", mockStore) + c.Set("repo", &model.Repo{ID: 1, FullName: "owner/repo"}) + LogStreamSSE(c) + }) + + ts = httptest.NewUnstartedServer(router) + ts.Config.WriteTimeout = writeTimeout + start(ts) + t.Cleanup(ts.Close) + + return ts, logService, done +} + +// TestLogStreamSSESurvivesWriteTimeout pins the rolling-deadline override for +// the log stream: a live LogStreamSSE stream keeps delivering log events well +// past the server-wide WriteTimeout, because the handler re-arms a rolling +// per-response write deadline (2*idlePingTime) before every write. It runs over +// both HTTP/1 and HTTP/2+TLS so the h2 responseWriter's SetWriteDeadline path +// is exercised too. Without the rolling deadline the server would tear the +// stream down at WriteTimeout and no log line would arrive after it — the +// assertion below would fail. +func TestLogStreamSSESurvivesWriteTimeout(t *testing.T) { + gin.SetMode(gin.TestMode) + + // The override proof needs the rolling deadline (2*idlePingTime) to outlast + // the server-wide WriteTimeout on the FIRST arm, so the stream survives past + // it even if the handler never gets to re-arm — otherwise survival would hinge + // on re-arming within the window every time, which a scheduler stall under + // load can miss (flaky). Keep 2*idlePingTime (200ms) > testWriteTimeout (80ms) + // with margin. Mutates the shared idlePingTime global, so no t.Parallel(). + origPing := idlePingTime + idlePingTime = 100 * time.Millisecond + t.Cleanup(func() { idlePingTime = origPing }) + + const testWriteTimeout = 80 * time.Millisecond + const stepID int64 = 42 + + transports := []struct { + name string + start func(*httptest.Server) + }{ + {"http1", func(ts *httptest.Server) { ts.Start() }}, + {"http2", func(ts *httptest.Server) { ts.EnableHTTP2 = true; ts.StartTLS() }}, + } + + for _, tr := range transports { + t.Run(tr.name, func(t *testing.T) { + ts, logService, handlerDone := newSSELogServer(t, testWriteTimeout, tr.start) + + ctx, cancel := context.WithTimeout(t.Context(), 5*time.Second) + t.Cleanup(cancel) + + req, err := http.NewRequestWithContext(ctx, http.MethodGet, ts.URL+"/stream/logs/1/1/42", nil) + require.NoError(t, err) + resp, err := ts.Client().Do(req) + require.NoError(t, err) + t.Cleanup(func() { _ = resp.Body.Close() }) + require.Equal(t, http.StatusOK, resp.StatusCode) + + // Emit log events continuously, past the WriteTimeout window. + pubStop := make(chan struct{}) + go func() { + ticker := time.NewTicker(30 * time.Millisecond) + defer ticker.Stop() + n := 0 + for { + select { + case <-pubStop: + return + case <-ticker.C: + n++ + _ = logService.Write(context.Background(), stepID, []*model.LogEntry{ + {Line: n, Data: []byte(fmt.Sprintf("log line %d", n))}, + }) + } + } + }() + t.Cleanup(func() { close(pubStop) }) + + // Read the stream and confirm at least one line arrives strictly + // after WriteTimeout has elapsed — proof the rolling deadline + // overrode the server-wide one. + start := time.Now() + br := bufio.NewReader(resp.Body) + readDeadline := time.Now().Add(2*testWriteTimeout + time.Second) + var gotLate bool + for time.Now().Before(readDeadline) { + line, err := br.ReadString('\n') + if err != nil { + break + } + if strings.TrimSpace(line) == "" { + continue + } + if time.Since(start) > testWriteTimeout { + gotLate = true + break + } + } + require.True(t, gotLate, + "log SSE stream stopped delivering before WriteTimeout elapsed; the rolling write deadline did not override the server-wide WriteTimeout") + + // Deterministically reap the handler goroutine before this subtest + // returns. Canceling the request trips the handler's + // requestCtx.Done() arm; waiting on handlerDone establishes a + // happens-before with the parent's idlePingTime restore, so the + // live handler can no longer race the restore's write of the + // global. This is the fix for that race — not a reason to rerun + // the test until it passes. + cancel() + select { + case <-handlerDone: + case <-time.After(2*testWriteTimeout + time.Second): + t.Fatal("log SSE handler did not return after request cancel; goroutine leaked") + } + }) + } +} + +// TestLogStreamSSEWedgeKill pins the return-on-write-error path for the log +// stream: when a client stops reading, the send buffer fills, and the blocked +// write is reclaimed by the rolling deadline (2*idlePingTime) — the handler +// observes the write error and returns instead of spinning forever. +// WriteTimeout is left 0 so the ONLY thing that can unblock the wedged write is +// the per-response deadline the fix arms; without it the handler would block +// indefinitely and this test would time out. +func TestLogStreamSSEWedgeKill(t *testing.T) { + gin.SetMode(gin.TestMode) + + // Mutates the shared idlePingTime global, so this test must NOT run + // t.Parallel(). + origPing := idlePingTime + idlePingTime = 40 * time.Millisecond + t.Cleanup(func() { idlePingTime = origPing }) + + const stepID int64 = 42 + + // WriteTimeout 0: no server-wide deadline. Only armSSEWriteDeadline can + // reclaim the wedged write. + ts, logService, handlerDone := newSSELogServer(t, 0, func(ts *httptest.Server) { ts.Start() }) + + ctx, cancel := context.WithCancelCause(t.Context()) + t.Cleanup(func() { cancel(nil) }) + + req, err := http.NewRequestWithContext(ctx, http.MethodGet, ts.URL+"/stream/logs/1/1/42", nil) + require.NoError(t, err) + resp, err := ts.Client().Do(req) + require.NoError(t, err) + t.Cleanup(func() { _ = resp.Body.Close() }) + require.Equal(t, http.StatusOK, resp.StatusCode) + + // The client deliberately never reads resp.Body. Flood large log entries so + // the kernel send buffer fills and the next handler write blocks. + pubStop := make(chan struct{}) + go func() { + ticker := time.NewTicker(5 * time.Millisecond) + defer ticker.Stop() + big := bytes.Repeat([]byte("x"), 256<<10) + n := 0 + for { + select { + case <-pubStop: + return + case <-ticker.C: + n++ + _ = logService.Write(context.Background(), stepID, []*model.LogEntry{ + {Line: n, Data: big}, + }) + } + } + }() + t.Cleanup(func() { close(pubStop) }) + + // The wedged write must be reclaimed by the rolling deadline and the + // handler must return within 2*idlePingTime + slack. + const slack = 3 * time.Second + select { + case <-handlerDone: + // handler returned — the write error was surfaced and acted on. + case <-time.After(2*idlePingTime + slack): + t.Fatalf("log SSE handler did not return within 2*idlePingTime+slack (%s); "+ + "the wedged write was never reclaimed / the write error was not acted on", 2*idlePingTime+slack) + } +} + +// flushErrRecorder is an http.ResponseWriter that reports a flush failure — +// standing in for the real *http.response once its write deadline has passed. +// It implements ONLY FlushError (no plain Flush), so a controller that reaches +// it must surface the error. +type flushErrRecorder struct { + http.ResponseWriter + flushErr error + flushes int +} + +func (w *flushErrRecorder) FlushError() error { + w.flushes++ + return w.flushErr +} + +// ginShapedWriter mimics the one property of gin's *responseWriter that matters +// here: it wraps another writer, exposes Unwrap, and implements the PLAIN +// http.Flusher (Flush() with no return value) — the arm that silently discards +// a flush error. +type ginShapedWriter struct { + http.ResponseWriter + flushes int +} + +func (w *ginShapedWriter) Flush() { w.flushes++ } +func (w *ginShapedWriter) Unwrap() http.ResponseWriter { return w.ResponseWriter } + +// TestNewResponseControllerSurfacesFlushError pins the reason the helper +// unwraps at all, and it is the DISCRIMINATING test for the flush path. +// +// Why a unit test rather than an end-to-end wedge: to wedge a real socket you +// must first write enough bytes to fill it, and those writes are themselves an +// error path that reclaims the stream — which is exactly why +// TestLogStreamSSEWedgeKill (256KB floods) passes with or without the fix and +// cannot see a swallowed flush. Every end-to-end variant inherits that hole, so +// the flush behavior is pinned directly. +// +// Gin's *responseWriter implements the plain http.Flusher, so +// http.NewResponseController(ginWriter).Flush() matches its `case Flusher` arm +// and returns nil WITHOUT unwrapping to the underlying writer's FlushError — +// net/http/responsecontroller.go orders the Flusher case before the rwUnwrapper +// case. Every `if err := rc.Flush()` check upstack is then dead code, and a +// ping-only wedged stream — whose only socket write happens inside that flush — +// lingers up to a further idlePingTime instead of being reclaimed within the +// 2*idlePingTime the design promises. +func TestNewResponseControllerSurfacesFlushError(t *testing.T) { + t.Parallel() + + newPair := func() (*flushErrRecorder, *ginShapedWriter) { + underlying := &flushErrRecorder{ + ResponseWriter: httptest.NewRecorder(), + flushErr: os.ErrDeadlineExceeded, + } + return underlying, &ginShapedWriter{ResponseWriter: underlying} + } + + // The fix: unwrapping past the gin-shaped writer reaches FlushError. + fixedUnderlying, fixedWriter := newPair() + fixedErr := newResponseController(fixedWriter).Flush() + + // The bug: building the controller straight on the gin-shaped writer lets + // its plain Flush() win the type switch and swallow the error. + buggyUnderlying, buggyWriter := newPair() + buggyErr := http.NewResponseController(buggyWriter).Flush() + + require.ErrorIs(t, fixedErr, os.ErrDeadlineExceeded, + "newResponseController must unwrap past gin so FlushError reaches the caller; "+ + "a nil here means every flush-error check in the SSE handlers is dead code") + require.Equal(t, 1, fixedUnderlying.flushes, "the underlying FlushError must actually run") + require.Zero(t, fixedWriter.flushes, "the gin-shaped Flush() must have been bypassed") + + require.NoError(t, buggyErr, + "a plain http.Flusher cannot report an error — if this ever fails, the premise "+ + "behind the unwrap has changed and it may be removable") + require.Zero(t, buggyUnderlying.flushes, "gin's Flush() must not have reached FlushError") + + // Assert the two paths DIFFER. Checking either arm alone would pass against + // a broken controller; the difference is the actual claim. + require.NotEqual(t, fixedErr == nil, buggyErr == nil, + "unwrapped and non-unwrapped flushes must disagree — if they agree, this test "+ + "no longer discriminates and the fix is unproven") +} + +// requireDeliveryThroughFinalInterval reads an SSE body for runFor and fails +// unless a line arrived inside the LAST 2*idlePingTime of that window. +// +// That final-interval check is what makes the caller a re-arm test rather than +// a first-arm test. The initial arm before the loop buys exactly one +// 2*idlePingTime of runway; a stream still delivering after several multiples +// of it can only be doing so because the handler re-armed inside the loop. +// +// Reading is bounded without a timer: when the deadline reclaims the response +// the server closes the connection and ReadString returns an error, so the +// loop ends on the event rather than on a guess. +func requireDeliveryThroughFinalInterval(t *testing.T, body io.Reader, runFor time.Duration) { + t.Helper() + + start := time.Now() + stopAfter := start.Add(runFor) + finalInterval := stopAfter.Add(-2 * idlePingTime) + + var lastLineAt time.Time + var lines int + br := bufio.NewReader(body) + for { + line, err := br.ReadString('\n') + if err != nil { + break + } + if strings.TrimSpace(line) == "" { + continue + } + lastLineAt = time.Now() + lines++ + if lastLineAt.After(stopAfter) { + break + } + } + + require.NotZero(t, lines, "the stream delivered nothing at all") + require.False(t, lastLineAt.Before(finalInterval), + "stream stopped delivering %s into a %s run, with the last line %s before the "+ + "final %s window; the handler armed the write deadline once before the loop "+ + "and never re-armed inside it", + lastLineAt.Sub(start), runFor, finalInterval.Sub(lastLineAt), 2*idlePingTime) +} + +// TestEventStreamSSEReArmsThroughoutStream pins the re-arm contract that +// TestEventStreamSSESurvivesWriteTimeout is too coarse to see. That test sizes +// the rolling deadline (2*idlePingTime) longer than the server-wide +// WriteTimeout on purpose, so the single arm before the loop already carries +// the stream past the window and deleting every in-loop arm still passes it. +// +// Here the run is five times the rolling deadline and the assertion lands in +// the final interval, which no single arm can reach: only a handler that +// re-arms before each write keeps the response alive that long. +// +// HTTP/1 only. The transport does not change which arms run — it only changes +// the SetWriteDeadline implementation underneath, which the coarse tests +// already exercise over both. +func TestEventStreamSSEReArmsThroughoutStream(t *testing.T) { + gin.SetMode(gin.TestMode) + + // Mutates the shared idlePingTime global, so no t.Parallel(). + origPing := idlePingTime + idlePingTime = 100 * time.Millisecond + t.Cleanup(func() { idlePingTime = origPing }) + + // Five rolling-deadline windows: long enough that a stream carried only by + // the first arm is dead four windows before the assertion. + runFor := 5 * 2 * idlePingTime + + // WriteTimeout 0 so the server-wide deadline cannot be what ends the + // stream: the per-response arm is the only deadline in play. + ts, broker, handlerDone := newSSEEventServer(t, 0, func(ts *httptest.Server) { ts.Start() }) + + ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second) + t.Cleanup(cancel) + + req, err := http.NewRequestWithContext(ctx, http.MethodGet, ts.URL+"/stream/events", nil) + require.NoError(t, err) + resp, err := ts.Client().Do(req) + require.NoError(t, err) + t.Cleanup(func() { _ = resp.Body.Close() }) + require.Equal(t, http.StatusOK, resp.StatusCode) + + // Publish faster than idlePingTime so the data arm, not the ping arm, is + // what drives the loop — the ping arm would re-arm on its own schedule and + // blur which call site is under test. + pubStop := make(chan struct{}) + go func() { + ticker := time.NewTicker(20 * time.Millisecond) + defer ticker.Stop() + topic := map[string]struct{}{pubsub.PublicTopic: {}} + for { + select { + case <-pubStop: + return + case <-ticker.C: + _ = broker.Publish(context.Background(), topic, + pubsub.Message{Data: []byte(`{"pipeline":1}`)}) + } + } + }() + t.Cleanup(func() { close(pubStop) }) + + requireDeliveryThroughFinalInterval(t, resp.Body, runFor) + + // Reap the handler before the parent restores idlePingTime, so the live + // handler's read of the global cannot race that write. + cancel() + select { + case <-handlerDone: + case <-time.After(5 * time.Second): + t.Fatal("SSE handler did not return after request cancel; goroutine leaked") + } +} + +// TestLogStreamSSEReArmsThroughoutStream is TestEventStreamSSEReArmsThroughoutStream +// for the log stream: same reasoning, same mutant, the other handler's loop. +func TestLogStreamSSEReArmsThroughoutStream(t *testing.T) { + gin.SetMode(gin.TestMode) + + origPing := idlePingTime + idlePingTime = 100 * time.Millisecond + t.Cleanup(func() { idlePingTime = origPing }) + + const stepID int64 = 42 + runFor := 5 * 2 * idlePingTime + + ts, logService, handlerDone := newSSELogServer(t, 0, func(ts *httptest.Server) { ts.Start() }) + + ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second) + t.Cleanup(cancel) + + req, err := http.NewRequestWithContext(ctx, http.MethodGet, ts.URL+"/stream/logs/1/1/42", nil) + require.NoError(t, err) + resp, err := ts.Client().Do(req) + require.NoError(t, err) + t.Cleanup(func() { _ = resp.Body.Close() }) + require.Equal(t, http.StatusOK, resp.StatusCode) + + pubStop := make(chan struct{}) + go func() { + ticker := time.NewTicker(20 * time.Millisecond) + defer ticker.Stop() + n := 0 + for { + select { + case <-pubStop: + return + case <-ticker.C: + n++ + _ = logService.Write(context.Background(), stepID, []*model.LogEntry{ + {Line: n, Data: []byte(fmt.Sprintf("log line %d", n))}, + }) + } + } + }() + t.Cleanup(func() { close(pubStop) }) + + requireDeliveryThroughFinalInterval(t, resp.Body, runFor) + + cancel() + select { + case <-handlerDone: + case <-time.After(5 * time.Second): + t.Fatal("log SSE handler did not return after request cancel; goroutine leaked") + } +} + +// TestLogStreamSSEEOFMarkerReachesHealthyClient pins that the eof marker +// actually reaches a client that never stopped reading. +// +// The eof write lives in the ctx.Done() arm, and every other write site arms +// the deadline immediately before writing. If that arm alone is missing, the +// last arm in force can be arbitrarily old by the time the logs end, because +// two things reset the ping countdown without re-arming: any select arm firing +// restarts time.After(idlePingTime), and the replay path re-arms only inside +// `id > last`. So a reconnect replaying more than 2*idlePingTime of backlog +// runs the loop hard with an already-expired deadline, and the eof write fails +// against a perfectly healthy socket. The client's EventSource then reconnects +// forever waiting for an end that was written and dropped. +// +// The reproduction is exactly that shape: Last-Event-ID far ahead of the live +// ids so the loop churns through the replay path, a burst longer than the +// rolling deadline, then the log service closes to trigger eof. WriteTimeout is +// 0 so nothing but the per-response arm governs the write. +func TestLogStreamSSEEOFMarkerReachesHealthyClient(t *testing.T) { + gin.SetMode(gin.TestMode) + + origPing := idlePingTime + idlePingTime = 150 * time.Millisecond + t.Cleanup(func() { idlePingTime = origPing }) + + const stepID int64 = 42 + // Three rolling-deadline windows of replay: the deadline armed before the + // loop is long expired by the time the logs end. + const replayFor = 900 * time.Millisecond + + ts, logService, handlerDone := newSSELogServer(t, 0, func(ts *httptest.Server) { ts.Start() }) + + ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second) + t.Cleanup(cancel) + + req, err := http.NewRequestWithContext(ctx, http.MethodGet, ts.URL+"/stream/logs/1/1/42", nil) + require.NoError(t, err) + // Far ahead of any id this stream will reach, so `id > last` is never true + // and the loop drives without passing a write site. + req.Header.Set("Last-Event-ID", "100000") + resp, err := ts.Client().Do(req) + require.NoError(t, err) + t.Cleanup(func() { _ = resp.Body.Close() }) + require.Equal(t, http.StatusOK, resp.StatusCode) + + // A healthy client: it reads continuously and never stops, so nothing about + // the peer can explain a failed write. + gotEOF := make(chan struct{}) + go func() { + br := bufio.NewReader(resp.Body) + for { + line, err := br.ReadString('\n') + if err != nil { + return + } + if strings.HasPrefix(line, "event: eof") { + close(gotEOF) + return + } + } + }() + + burstDone := time.After(replayFor) + ticker := time.NewTicker(5 * time.Millisecond) + defer ticker.Stop() + for n := 0; ; { + stop := false + select { + case <-burstDone: + stop = true + case <-ticker.C: + n++ + require.NoError(t, logService.Write(context.Background(), stepID, + []*model.LogEntry{{Line: n, Data: []byte(fmt.Sprintf("replayed line %d", n))}})) + } + if stop { + break + } + } + + // Closing the stream ends Tail, which cancels the handler's context with + // context.Canceled — the eof arm. + require.NoError(t, logService.Close(context.Background(), stepID)) + + select { + case <-gotEOF: + case <-handlerDone: + // The handler returned without the marker ever landing: the eof write + // or its flush failed on an expired deadline. + select { + case <-gotEOF: + case <-time.After(time.Second): + t.Fatal("log stream ended without delivering the eof marker to a client that " + + "never stopped reading; the eof write ran against a write deadline that " + + "expired during replay and was never re-armed") + } + case <-time.After(10 * time.Second): + t.Fatal("neither the eof marker nor the handler's return arrived") + } +} + +// A stream carrying only keep-alive pings is the case the absolute ceiling +// exists for, and the one the rolling deadline cannot reach: a 9-byte ping +// never fills a socket send buffer, so its write never blocks and no write +// deadline it arms can ever expire. Left to the rolling deadline alone such a +// stream is held for as long as the peer keeps the connection open without +// reading — measured at 75x the interval the spec promises, and unbounded in +// principle. +// +// These pin the timer specifically. The clamp in armSSEWriteDeadline is covered +// by TestArmSSEWriteDeadlineClampsToTheCeiling, but the clamp is not what ends +// the stream: with the timer neutered and the clamp intact, a ping-only stream +// is not reclaimed at all. Only an end-to-end test distinguishes them. +func TestEventStreamSSEPingOnlyStreamEndsAtTheCeiling(t *testing.T) { + gin.SetMode(gin.TestMode) + shortenStreamBoundsForTest(t) + + // WriteTimeout 0: with no server-wide deadline in play, the ceiling is the + // only thing that can end this handler. + ts, _, handlerDone := newSSEEventServer(t, 0, func(ts *httptest.Server) { ts.Start() }) + + resp := getAndReadFirstPing(t, ts.URL+"/stream/events") + // Stop reading without closing, and publish nothing: pings only, to a peer + // that never drains them. + t.Cleanup(func() { _ = resp.Body.Close() }) + + select { + case <-handlerDone: + case <-time.After(streamMaxDuration + 5*time.Second): + t.Fatal("ping-only stream outlived its absolute ceiling; the handler was never reclaimed") + } +} + +func TestLogStreamSSEPingOnlyStreamEndsAtTheCeiling(t *testing.T) { + gin.SetMode(gin.TestMode) + shortenStreamBoundsForTest(t) + + ts, _, handlerDone := newSSELogServer(t, 0, func(ts *httptest.Server) { ts.Start() }) + + resp := getAndReadFirstPing(t, ts.URL+"/stream/logs/1/1/42") + t.Cleanup(func() { _ = resp.Body.Close() }) + + select { + case <-handlerDone: + case <-time.After(streamMaxDuration + 5*time.Second): + t.Fatal("ping-only log stream outlived its absolute ceiling; the handler was never reclaimed") + } +} + +// shortenStreamBoundsForTest scales the ping interval and the ceiling down so a +// ceiling that is an hour in production is observable in a test. Mutating the +// package globals means these tests must not run with t.Parallel(), per the +// rule documented on idlePingTime. +func shortenStreamBoundsForTest(t *testing.T) { + t.Helper() + origPing, origMax := idlePingTime, streamMaxDuration + idlePingTime, streamMaxDuration = 50*time.Millisecond, 500*time.Millisecond + t.Cleanup(func() { idlePingTime, streamMaxDuration = origPing, origMax }) +} + +// getAndReadFirstPing opens the stream and consumes its opening ping, so the +// caller knows the handler has reached its loop before it stops reading. +func getAndReadFirstPing(t *testing.T, url string) *http.Response { + t.Helper() + + req, err := http.NewRequestWithContext(t.Context(), http.MethodGet, url, nil) + require.NoError(t, err) + resp, err := http.DefaultClient.Do(req) + require.NoError(t, err) + + buf := make([]byte, len(": ping\n\n")) + _, err = io.ReadFull(resp.Body, buf) + require.NoError(t, err, "the stream should open with a ping") + require.Equal(t, ": ping\n\n", string(buf)) + + return resp +} diff --git a/server/api/user.go b/server/api/user.go index c34213411e0..06f327bb9b5 100644 --- a/server/api/user.go +++ b/server/api/user.go @@ -15,7 +15,9 @@ package api import ( + "context" "encoding/base32" + "errors" "net/http" "strconv" "strings" @@ -123,13 +125,29 @@ func GetRepos(c *gin.Context) { dbStaleReposMap[r.ID] = r } + // Same slow shape as the permission sync below: paging against the forge + // can outrun the server-wide WriteTimeout, which bounds total handler + // runtime. This is the "add repository" listing, so a large account is + // the expected case rather than the edge one. + onPage := slowHandlerProgress(c, newResponseController(c.Writer)) _repos, err := utils.Paginate(func(page int) ([]*model.Repo, error) { + if err := onPage(); err != nil { + return nil, err + } return _forge.Repos(c, user, &model.ListOptions{ Page: page, PerPage: perPage, }) }, maxPage) if err != nil { + // The error may be onPage's rather than the forge's: the client went + // away mid-paging, so there is no listing to return and no fault to + // report. + if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) { + log.Debug().Err(err).Msgf("get repos: client has gone away while listing repositories for user %s", user.Login) + c.AbortWithStatus(statusClientClosedRequest) + return + } c.String(http.StatusInternalServerError, "Error fetching repository list. %s", err) return } @@ -228,7 +246,20 @@ func RefreshRepos(c *gin.Context) { return } - if err := updateRepoPermissions(c, user, _store, _forge, user.ForgeID); err != nil { + // The sync pages against the forge and can outrun the server-wide + // WriteTimeout, which bounds total handler runtime; re-arm per page so a + // refresh over a large account still returns its 200, and stop paging once + // the client has gone away. + onPage := slowHandlerProgress(c, newResponseController(c.Writer)) + if err := updateRepoPermissions(c, user, _store, _forge, user.ForgeID, onPage); err != nil { + // onPage stopped the sync because the client went away, not because the + // sync failed. A 500 here would report our own healthy refresh as a + // server fault. + if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) { + log.Debug().Err(err).Msgf("refresh repos: client has gone away while syncing repo permissions for user %s", user.Login) + c.AbortWithStatus(statusClientClosedRequest) + return + } log.Error().Err(err).Msgf("Can't update repo permissions for user %s in forge %s", user.Login, _forge.Name()) c.AbortWithStatus(http.StatusInternalServerError) return diff --git a/server/api/users_test.go b/server/api/users_test.go index 2d319a55ec1..922513cd1d7 100644 --- a/server/api/users_test.go +++ b/server/api/users_test.go @@ -17,13 +17,19 @@ package api import ( + "context" "net/http" + "net/http/httptest" "testing" "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/mock" "github.com/stretchr/testify/require" + "go.woodpecker-ci.org/woodpecker/v3/server" + forge_mocks "go.woodpecker-ci.org/woodpecker/v3/server/forge/mocks" "go.woodpecker-ci.org/woodpecker/v3/server/model" + manager_mocks "go.woodpecker-ci.org/woodpecker/v3/server/services/mocks" ) func TestPostUser(t *testing.T) { @@ -60,3 +66,88 @@ func TestPostUser(t *testing.T) { assert.EqualValues(t, 2, created.ForgeID) }) } + +// installUserForge wires a mock manager whose ForgeFromUser returns forge, so +// GetRepos and RefreshRepos resolve the forge they page against. +func installUserForge(t *testing.T) *forge_mocks.MockForge { + t.Helper() + mgr := manager_mocks.NewMockManager(t) + _forge := forge_mocks.NewMockForge(t) + mgr.On("ForgeFromUser", mock.Anything).Return(_forge, nil) + server.Config.Services.Manager = mgr + return _forge +} + +// The slow forge-paging handlers gained a 499-vs-500 split: a client that +// disconnects mid-paging is reported as 499 (client closed request), a genuine +// forge error as 500. The slowHandlerProgress hook checks the request context +// before each page, so a request whose context is already canceled makes the +// forge call return context.Canceled without any real disconnect timing — the +// same error the abort branch keys on. +func TestGetReposClassifiesClientCancelVsForgeError(t *testing.T) { + s := newTestStore(t) + user := &model.User{ID: 1, ForgeID: defaultForgeID, Login: "alice"} + + t.Run("client cancel while listing repos returns 499", func(t *testing.T) { + installUserForge(t) + tc := newTestContext(t, s) + withUser(user)(tc) + + // all=true takes the forge-paging path; a canceled request context + // trips slowHandlerProgress before the first forge page. + ctx, cancel := context.WithCancelCause(t.Context()) + cancel(nil) + tc.Ctx.Request = httptest.NewRequestWithContext(ctx, http.MethodGet, "/user/repos?all=true", nil) + + GetRepos(tc.Ctx) + + assert.Equal(t, statusClientClosedRequest, tc.Recorder.Code, tc.Recorder.Body.String()) + }) + + t.Run("forge error while listing repos returns 500", func(t *testing.T) { + _forge := installUserForge(t) + _forge.On("Repos", mock.Anything, mock.Anything, mock.Anything). + Return(nil, assert.AnError) + tc := newTestContext(t, s) + withUser(user)(tc) + tc.Ctx.Request = httptest.NewRequest(http.MethodGet, "/user/repos?all=true", nil) + + GetRepos(tc.Ctx) + + assert.Equal(t, http.StatusInternalServerError, tc.Recorder.Code, tc.Recorder.Body.String()) + }) +} + +func TestRefreshReposClassifiesClientCancelVsForgeError(t *testing.T) { + s := newTestStore(t) + user := &model.User{ID: 1, ForgeID: defaultForgeID, Login: "alice"} + + t.Run("client cancel while syncing permissions returns 499", func(t *testing.T) { + installUserForge(t) + tc := newTestContext(t, s) + withUser(user)(tc) + + ctx, cancel := context.WithCancelCause(t.Context()) + cancel(nil) + tc.Ctx.Request = httptest.NewRequestWithContext(ctx, http.MethodGet, "/user/repos", nil) + + RefreshRepos(tc.Ctx) + + assert.Equal(t, statusClientClosedRequest, tc.Recorder.Code, tc.Recorder.Body.String()) + }) + + t.Run("forge error while syncing permissions returns 500", func(t *testing.T) { + _forge := installUserForge(t) + _forge.On("Repos", mock.Anything, mock.Anything, mock.Anything). + Return(nil, assert.AnError) + // RefreshRepos logs the forge name on the 500 path. + _forge.On("Name").Return("mock-forge").Maybe() + tc := newTestContext(t, s) + withUser(user)(tc) + tc.Ctx.Request = httptest.NewRequest(http.MethodGet, "/user/repos", nil) + + RefreshRepos(tc.Ctx) + + assert.Equal(t, http.StatusInternalServerError, tc.Recorder.Code, tc.Recorder.Body.String()) + }) +} diff --git a/server/api/writedeadline.go b/server/api/writedeadline.go new file mode 100644 index 00000000000..3748674c18a --- /dev/null +++ b/server/api/writedeadline.go @@ -0,0 +1,128 @@ +// Copyright 2026 Woodpecker Authors +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package api + +import ( + "net/http" + "sync" + "time" + + "github.com/gin-gonic/gin" + "github.com/rs/zerolog/log" +) + +// The server-wide WriteTimeout is a wall-clock budget on the whole handler, not +// only on a blocked write: net/http arms it once from the request-header read +// (net/http/server.go, the deferred SetWriteDeadline in readRequest), so handler +// compute counts against it too. That is what bounds a write wedged on a +// zero-window peer, but it also caps every legitimately slow response. +// +// A handler that can outrun the budget takes a rolling per-response deadline +// instead, refreshing it as it makes progress. The helpers here are that seam, +// shared by the SSE streams and by the handlers that page against a forge. + +// newResponseController builds the ResponseController a handler drives, +// unwrapping past gin's *responseWriter deliberately. +// +// Gin implements the plain http.Flusher (Flush() with no return), so +// ResponseController.Flush() matches its `case Flusher` arm and returns nil +// without ever reaching the wrapped *http.response.FlushError() — every +// flush-error check upstack would be dead code. Unwrapping restores FlushError. +// Gin's rw.Write has already run WriteHeaderNow by then, so header bookkeeping +// and SetWriteDeadline are unaffected. SetWriteDeadline alone would not need +// this, since ResponseController walks Unwrap itself; Flush is what forces it. +// +// The unwrap loops rather than stepping once. One step suffices today, because +// gin's writer is outermost — but a middleware that wraps it (W → gin → +// *http.response) would leave a single step on gin's own writer, whose plain +// Flush() wins the Flusher arm and silently restores the swallowed-flush bug +// this helper exists to prevent. Stopping at the first writer that reports +// flush errors keeps that fix independent of how deep the writer is buried. +// +// A writer with no FlushError anywhere in its chain is used as-is: degraded, +// not broken. +func newResponseController(rw http.ResponseWriter) *http.ResponseController { + for { + if _, ok := rw.(interface{ FlushError() error }); ok { + break + } + u, ok := rw.(interface{ Unwrap() http.ResponseWriter }) + if !ok { + break + } + rw = u.Unwrap() + } + return http.NewResponseController(rw) +} + +// extendWriteDeadline pushes the response's write deadline out to now+d, +// replacing whatever the server-wide WriteTimeout armed. It is the primitive +// behind both rolling-deadline shapes: a stream re-arms before each write, a +// slow handler re-arms as it makes progress. +// +// A writer that does not support SetWriteDeadline is reported at Warn, not +// Debug: it means the per-response bound is not in force and the handler is +// running under the server-wide deadline alone, which is the condition this +// whole mechanism exists to avoid. It is not fatal — the server-wide deadline +// still reclaims the connection — but an operator should not need debug logging +// to discover the mitigation is inactive. +// +// The warnOnce keeps that Warn to one line per response, and is required rather +// than optional because every caller arms repeatedly against one response: a +// stream before each write, a slow handler on each unit of progress, the queue +// long-poll at 10Hz. The condition is a property of the writer and cannot +// change mid-response, so repeating the line adds no information and floods +// the log for as long as the request lives. +func extendWriteDeadline(rc *http.ResponseController, d time.Duration, warnOnce *sync.Once) { + err := rc.SetWriteDeadline(time.Now().Add(d)) + if err == nil { + return + } + warnOnce.Do(func() { + log.Warn().Err(err).Dur("extension", d). + Msg("SetWriteDeadline unsupported: no per-response write bound, " + + "handler runs under the server-wide WriteTimeout alone") + }) +} + +// slowHandlerProgress returns the per-iteration hook a slow handler runs as it +// makes progress: it re-arms the rolling write deadline and reports whether the +// client is still there. +// +// Both halves are required together, and that is the point of pairing them. +// Re-arming removes the server-wide WriteTimeout as a bound, and a handler that +// pages against a forge writes nothing while it works — so no write can trip +// the deadline and nothing else would ever end the loop. Cancellation is the +// bound that replaces the one the arming removed. +func slowHandlerProgress(c *gin.Context, rc *http.ResponseController) func() error { + requestCtx := c.Request.Context() + var warnOnce sync.Once + return func() error { + if err := requestCtx.Err(); err != nil { + return err + } + extendWriteDeadline(rc, slowHandlerWriteExtension, &warnOnce) + return nil + } +} + +// statusClientClosedRequest is nginx's 499: the client closed the connection +// before the server produced a response. Not an IANA code, so net/http has no +// constant for it, but gin writes it verbatim and it is what the access log +// needs to distinguish an abandoned request from a completed one. A handler +// whose slowHandlerProgress hook reports cancellation has done partial work +// nobody will read; without an explicit status gin finalizes a bare 200 and +// the log records that partial work as a success. +const statusClientClosedRequest = 499 diff --git a/server/api/writedeadline_test.go b/server/api/writedeadline_test.go new file mode 100644 index 00000000000..55de8eb2fd0 --- /dev/null +++ b/server/api/writedeadline_test.go @@ -0,0 +1,209 @@ +// Copyright 2026 Woodpecker Authors +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package api + +import ( + "context" + "net/http" + "net/http/httptest" + "sync" + "testing" + "time" + + "github.com/gin-gonic/gin" + "github.com/stretchr/testify/assert" +) + +// deadlineRecorder is a ResponseWriter that records SetWriteDeadline calls. +// It deliberately does NOT implement Flush, so a controller built on it +// reports whether the unwrap reached this writer or stopped at a wrapper. +type deadlineRecorder struct { + http.ResponseWriter + deadlines []time.Time + err error +} + +func (d *deadlineRecorder) SetWriteDeadline(t time.Time) error { + if d.err != nil { + return d.err + } + d.deadlines = append(d.deadlines, t) + return nil +} + +// flushWrapper mimics gin's *responseWriter: it implements the plain +// http.Flusher (Flush with no error return) and exposes the wrapped writer +// through Unwrap. A ResponseController built on this rather than on its +// unwrapped target matches the Flusher arm and silently swallows flush errors. +type flushWrapper struct { + http.ResponseWriter + flushed bool +} + +func (f *flushWrapper) Flush() { f.flushed = true } +func (f *flushWrapper) Unwrap() http.ResponseWriter { return f.ResponseWriter } + +func TestNewResponseControllerUnwrapsPastAFlusher(t *testing.T) { + // The wrapper implements Flush; the target does not. If the controller is + // built on the wrapper it matches the Flusher arm and returns nil. Built on + // the unwrapped target it finds neither FlushError nor Flusher and reports + // ErrNotSupported. That difference is the whole reason for the unwrap, so + // asserting it is what keeps the unwrap from being quietly removed. + target := &deadlineRecorder{} + wrapper := &flushWrapper{ResponseWriter: target} + + err := newResponseController(wrapper).Flush() + assert.ErrorIs(t, err, http.ErrNotSupported, + "controller must drive the unwrapped writer, not the wrapper's own Flush") + assert.False(t, wrapper.flushed, + "the wrapper's Flush must not be what runs") + + // Built on the wrapper directly, the swallowed-flush behavior appears — + // this is the bug the unwrap avoids, pinned so the two paths must differ. + assert.NoError(t, http.NewResponseController(wrapper).Flush()) + assert.True(t, wrapper.flushed) +} + +func TestNewResponseControllerHandlesAWriterWithoutUnwrap(t *testing.T) { + // Degraded, not broken: a writer that cannot be unwrapped is used as-is. + target := &deadlineRecorder{} + + before := time.Now() + extendWriteDeadline(newResponseController(target), time.Minute, &sync.Once{}) + + if assert.Len(t, target.deadlines, 1) { + assert.WithinRange(t, target.deadlines[0], + before.Add(time.Minute), time.Now().Add(time.Minute)) + } +} + +func TestExtendWriteDeadlineSurvivesAnUnsupportedWriter(t *testing.T) { + // A writer that cannot take a deadline must not panic or block the caller: + // the handler keeps running under the server-wide deadline instead. + target := &deadlineRecorder{err: http.ErrNotSupported} + + assert.NotPanics(t, func() { + extendWriteDeadline(newResponseController(target), time.Minute, &sync.Once{}) + }) + assert.Empty(t, target.deadlines) +} + +func TestArmSSEWriteDeadlineClampsToTheCeiling(t *testing.T) { + // The ceiling is what reclaims a ping-only stream, whose writes never block + // and so never trip the rolling deadline on their own. + for _, tc := range []struct { + name string + remaining time.Duration + wantMax time.Duration + }{ + {"far from the ceiling: the full rolling arm", time.Hour, 2 * idlePingTime}, + {"near the ceiling: clamped to what remains", idlePingTime / 2, idlePingTime}, + {"past the ceiling: floored, never in the past", -time.Minute, idlePingTime}, + } { + t.Run(tc.name, func(t *testing.T) { + target := &deadlineRecorder{} + before := time.Now() + + armSSEWriteDeadline(newResponseController(target), before.Add(tc.remaining), &sync.Once{}) + + if assert.Len(t, target.deadlines, 1) { + got := target.deadlines[0] + assert.False(t, got.Before(before), + "an arm must never set a deadline in the past: %v", got) + assert.LessOrEqual(t, got.Sub(before), tc.wantMax+time.Second, + "arm must not exceed the clamp") + } + }) + } +} + +// newSlowHandlerCtx builds a gin context whose request carries the given +// context, plus a ResponseController over a deadlineRecorder, mirroring the +// construction the other tests here use (gin.CreateTestContext + +// newResponseController over a deadline-recording writer). +func newSlowHandlerCtx(reqCtx context.Context) (*gin.Context, *http.ResponseController, *deadlineRecorder) { + target := &deadlineRecorder{} + c, _ := gin.CreateTestContext(httptest.NewRecorder()) + c.Request = httptest.NewRequest(http.MethodGet, "/", nil).WithContext(reqCtx) + return c, newResponseController(target), target +} + +func TestSlowHandlerProgressContinuesWhileTheClientIsThere(t *testing.T) { + // The live path: a still-connected client means the hook re-arms the + // rolling deadline and returns nil, so the paging loop keeps going. This is + // the branch the four handlers do NOT translate to 499 — a false positive + // here would abort a healthy request as a client disconnect. + c, rc, target := newSlowHandlerCtx(context.Background()) + + hook := slowHandlerProgress(c, rc) + + before := time.Now() + assert.NoError(t, hook(), "a live request context must not end the paging loop") + if assert.Len(t, target.deadlines, 1, "the hook must re-arm the write deadline each call") { + assert.WithinRange(t, target.deadlines[0], + before.Add(slowHandlerWriteExtension), time.Now().Add(slowHandlerWriteExtension)) + } + + // Each call re-arms afresh: a slow handler leans on this per-iteration. + assert.NoError(t, hook()) + assert.Len(t, target.deadlines, 2) +} + +func TestSlowHandlerProgressStopsOnCancellation(t *testing.T) { + // The cancel path: once the request context reports an error the hook + // returns it verbatim, and the error must be Is-matchable so the handlers + // can distinguish a client disconnect (→499, stop paging) from a genuine + // forge failure (→500). The hook must also NOT re-arm the deadline once + // canceled — it bails before the extend. + for _, tc := range []struct { + name string + arm func() context.Context + wantErr error + }{ + { + "client hang-up: context.Canceled", + func() context.Context { + ctx, cancel := context.WithCancelCause(context.Background()) + cancel(nil) + return ctx + }, + context.Canceled, + }, + { + "budget blown: context.DeadlineExceeded", + func() context.Context { + ctx, cancel := context.WithDeadline(context.Background(), time.Now().Add(-time.Minute)) + t.Cleanup(cancel) + return ctx + }, + context.DeadlineExceeded, + }, + } { + t.Run(tc.name, func(t *testing.T) { + // The context returned by arm is already terminal (canceled now, or + // past its deadline by construction), so no sleeps are needed. + reqCtx := tc.arm() + + c, rc, target := newSlowHandlerCtx(reqCtx) + err := slowHandlerProgress(c, rc)() + + assert.Error(t, err, "a canceled request context must end the paging loop") + assert.ErrorIs(t, err, tc.wantErr, + "the error must be Is-matchable so handlers can translate it to 499") + assert.Empty(t, target.deadlines, + "a canceled hook must bail before re-arming the write deadline") + }) + } +}