Compare commits
5
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
0945831442 | ||
|
|
f82b730c31 | ||
|
|
1a1fee0874 | ||
|
|
290925f184 | ||
|
|
73353bc8e5 |
@@ -3245,9 +3245,9 @@ each hook. The order, read off the fx stop-hook log:
|
||||
|
||||
1. `ArchiveSweeper`
|
||||
2. `RetentionReaper`
|
||||
3. `server` — the HTTP drain, bounded separately by
|
||||
`server.ShutdownTimeout` (**3 seconds**), then a Sentry flush if
|
||||
`SENTRY_DSN` is set
|
||||
3. `server` — the HTTP drain, bounded by `server.ShutdownTimeout`
|
||||
(**3 seconds**) and by what the hooks before it left, then a Sentry
|
||||
flush if `SENTRY_DSN` is set
|
||||
4. `delivery.Engine` — waits for its workers, then closes the archive
|
||||
databases
|
||||
5. `healthcheck`
|
||||
@@ -3267,23 +3267,30 @@ exhaust the sequence budget at the instant it finished, and every
|
||||
later hook — the delivery engine, the healthcheck, the webhook DB
|
||||
manager and the database close — would be skipped in exactly the
|
||||
case where the drain mattered. 3 seconds leaves 2 seconds
|
||||
(`server.TailHookReserve`) for the tail, which is far more than the
|
||||
microseconds it needs.
|
||||
(`server.TailHookReserve`) for the tail. The reserve is that
|
||||
remainder, not a figure sized to the tail, which takes about a
|
||||
millisecond.
|
||||
|
||||
That reserve belongs to the tail hooks, not to the server hook, and
|
||||
the Sentry flush is what could take it: it runs after the drain
|
||||
**inside the same hook**, and `sentry.Flush` takes a bare duration
|
||||
and honours no context, so an unreachable Sentry endpoint would add
|
||||
its own timeout on top of a full-length drain and consume the whole
|
||||
sequence budget by itself. It is therefore clamped to whatever is
|
||||
left on the stop context minus the reserve, and skipped when that
|
||||
leaves too little to be worth attempting — so a full-length drain
|
||||
means Sentry events are dropped rather than the database close being
|
||||
skipped.
|
||||
the server hook could take it in two ways. The hooks before it may
|
||||
already have spent part of the budget, so a full 3-second drain
|
||||
would come out of the reserve; the drain is therefore also bounded
|
||||
by whatever is left on the stop context minus the reserve. And the
|
||||
Sentry flush runs after the drain **inside the same hook**, and
|
||||
`sentry.Flush` takes a bare duration and honours no context, so an
|
||||
unreachable Sentry endpoint would add its own timeout on top of a
|
||||
full-length drain and consume the whole sequence budget by itself.
|
||||
It is clamped the same way, and skipped when that leaves too little
|
||||
to be worth attempting — so a full-length drain means Sentry events
|
||||
are dropped rather than the database close being skipped.
|
||||
|
||||
This does not make the database close unconditional: a wedged
|
||||
`ArchiveSweeper` or `RetentionReaper` still runs first and can
|
||||
consume the whole budget on its own.
|
||||
This does not make the database close unconditional. A slow
|
||||
`ArchiveSweeper` or `RetentionReaper` is enough to cut the shutdown
|
||||
short, not only one that consumes the whole budget: what they spend
|
||||
comes out of the drain first, so after 2 seconds of theirs a request
|
||||
still in flight gets 1 second to finish, and after 3 it gets none.
|
||||
Past 3 seconds they spend the reserve itself, and one that takes the
|
||||
whole budget skips every hook after it, the database close included.
|
||||
|
||||
The value is chosen to sit inside the container stop grace period.
|
||||
Docker's default `docker stop` grace is 10 seconds and the Dockerfile
|
||||
|
||||
@@ -38,17 +38,19 @@ import (
|
||||
// hook that used the whole budget would exhaust it at that instant,
|
||||
// and fx would skip every hook after the server — the delivery
|
||||
// engine, the healthcheck, the webhook DB manager and the database
|
||||
// close. That hook is the 3s HTTP drain plus the Sentry flush that
|
||||
// follows it in the same hook, so the flush is clamped to the stop
|
||||
// close. That hook is the HTTP drain plus the Sentry flush that
|
||||
// follows it in the same hook, and each is clamped to the stop
|
||||
// context's remaining time less server.TailHookReserve rather than
|
||||
// running for its own fixed 2s; the reserve is what the tail hooks
|
||||
// live on, and they are microsecond-scale in normal operation.
|
||||
// running for its own fixed 3s and 2s; the reserve is what the tail
|
||||
// hooks live on, and they are microsecond-scale in normal operation.
|
||||
// TestStopTimeout_LeavesHeadroomForTailHooks pins the arithmetic
|
||||
// across every drain length.
|
||||
// across every drain length and every amount of budget the hooks
|
||||
// before the server may already have spent.
|
||||
//
|
||||
// This does not make the database close unconditional: the
|
||||
// ArchiveSweeper and RetentionReaper hooks run before the server
|
||||
// and can still consume the whole budget on their own.
|
||||
// ArchiveSweeper and RetentionReaper hooks run before the server.
|
||||
// What they spend comes out of the drain first, but past 3s it comes
|
||||
// out of the reserve, and they can consume the whole budget.
|
||||
const stopTimeout = 5 * time.Second
|
||||
|
||||
// exitUsage is the status for a command line this binary cannot make
|
||||
|
||||
@@ -252,22 +252,40 @@ const tailHeadroom = 2 * time.Second
|
||||
// can produce, since a shorter drain leaves the flush more room and
|
||||
// the worst case is not necessarily at either extreme.
|
||||
//
|
||||
// Shrinking either budget, or unbounding the flush again, must fail
|
||||
// here rather than silently recreating a hook that swallows the
|
||||
// whole sequence.
|
||||
// Nor does the hook start on a full budget: the ArchiveSweeper and
|
||||
// RetentionReaper hooks run before it, and whatever they spent is
|
||||
// gone. The outer sweep walks every amount they can spend. Once they
|
||||
// have eaten into the headroom themselves, the hook must spend
|
||||
// nothing of what is left. A drain that starts on the full budget
|
||||
// must still get all of ShutdownTimeout, so a smaller stopTimeout
|
||||
// cannot silently shorten every drain.
|
||||
//
|
||||
// Shrinking either budget, or unbounding the drain or the flush
|
||||
// again, must fail here rather than silently recreating a hook that
|
||||
// swallows the whole sequence.
|
||||
func TestStopTimeout_LeavesHeadroomForTailHooks(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
require.Less(t, server.ShutdownTimeout, stopTimeout)
|
||||
require.Equal(
|
||||
t, server.ShutdownTimeout, server.DrainBudget(stopTimeout),
|
||||
"a drain that starts on the full stop budget is cut short",
|
||||
)
|
||||
|
||||
const step = 10 * time.Millisecond
|
||||
|
||||
for drain := time.Duration(0); drain <= server.ShutdownTimeout; drain += step {
|
||||
hook := drain + server.SentryFlushBudget(stopTimeout-drain)
|
||||
for spent := time.Duration(0); spent <= stopTimeout; spent += step {
|
||||
remaining := stopTimeout - spent
|
||||
longest := max(server.DrainBudget(remaining), 0)
|
||||
|
||||
require.LessOrEqual(
|
||||
t, hook+tailHeadroom, stopTimeout,
|
||||
"a %s drain leaves the tail hooks short", drain,
|
||||
)
|
||||
for drain := time.Duration(0); drain <= longest; drain += step {
|
||||
hook := drain + server.SentryFlushBudget(remaining-drain)
|
||||
|
||||
require.GreaterOrEqual(
|
||||
t, remaining-hook, min(remaining, tailHeadroom),
|
||||
"a %s drain after %s of earlier hooks leaves "+
|
||||
"the tail hooks short", drain, spent,
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1425,6 +1425,32 @@ func TestDeliverHTTP_InvalidConfig(t *testing.T) {
|
||||
)
|
||||
}
|
||||
|
||||
// TestDeliverHTTP_InvalidConfigUnrecordedStaysPending: a delivery is
|
||||
// failed for an invalid config only once the reason is recorded.
|
||||
// Unrecorded, it stays pending, where the sweep finds it again.
|
||||
func TestDeliverHTTP_InvalidConfigUnrecordedStaysPending(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
db := testWebhookDB(t)
|
||||
e := testEngine(t, 1)
|
||||
|
||||
event, del := iSeedEventAndDelivery(
|
||||
t, db, `{"config":"invalid"}`, "",
|
||||
)
|
||||
|
||||
task, d := iHTTPTaskAndDelivery(
|
||||
event, del, "bad-config", `not-json`, 0, 1,
|
||||
)
|
||||
|
||||
require.NoError(t, db.Exec("drop table delivery_results").Error)
|
||||
|
||||
e.ExportDeliverHTTP(context.TODO(), db, d, task)
|
||||
|
||||
iAssertStatus(t, db, del.ID,
|
||||
database.DeliveryStatusPending,
|
||||
)
|
||||
}
|
||||
|
||||
// --- Notify batching ---
|
||||
|
||||
func TestNotify_MultipleTasks(t *testing.T) {
|
||||
|
||||
@@ -5,6 +5,7 @@ import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
@@ -1056,6 +1057,21 @@ func TestParseHTTPConfig_MissingURL(t *testing.T) {
|
||||
)
|
||||
}
|
||||
|
||||
func TestParseHTTPConfig_Undecodable(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
e := testEngine(t, 1)
|
||||
|
||||
_, err := e.ExportParseHTTPConfig(
|
||||
`{"url":"https://example.com/hook","timeout":"soon"}`,
|
||||
)
|
||||
|
||||
assert.Error(t, err,
|
||||
"config that does not decode should return error, "+
|
||||
"even when the part that did names a URL",
|
||||
)
|
||||
}
|
||||
|
||||
func TestScheduleRetry_SendsToRetryChannel(
|
||||
t *testing.T,
|
||||
) {
|
||||
@@ -1241,6 +1257,33 @@ func TestDoHTTPRequest_ForwardsHeaders(t *testing.T) {
|
||||
)
|
||||
}
|
||||
|
||||
// A response that ends before the length it announced is an error, not
|
||||
// a short body.
|
||||
func TestDoHTTPRequest_CutShortResponseIsAnError(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
ts := httptest.NewServer(
|
||||
http.HandlerFunc(
|
||||
func(w http.ResponseWriter, _ *http.Request) {
|
||||
w.Header().Set("Content-Length", "100")
|
||||
_, _ = w.Write([]byte("cut short"))
|
||||
},
|
||||
),
|
||||
)
|
||||
defer ts.Close()
|
||||
|
||||
e := testEngine(t, 1)
|
||||
|
||||
_, body, _, err := e.ExportDoHTTPRequest(
|
||||
context.TODO(),
|
||||
&delivery.HTTPTargetConfig{URL: ts.URL},
|
||||
&database.Event{},
|
||||
)
|
||||
|
||||
require.ErrorIs(t, err, io.ErrUnexpectedEOF)
|
||||
assert.Empty(t, body)
|
||||
}
|
||||
|
||||
// The event's stored inbound headers carry the same Content-Type the
|
||||
// receiver saved as the event's ContentType, so a delivery could send
|
||||
// it twice. It must go out exactly once, with a Content-Type configured
|
||||
@@ -1317,6 +1360,34 @@ func TestApplyRequestHeaders_SendsOneContentType(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// Stored inbound headers that do not decode forward nothing, not the
|
||||
// part of them that happened to decode.
|
||||
func TestApplyRequestHeaders_UndecodableInboundForwardsNothing(
|
||||
t *testing.T,
|
||||
) {
|
||||
t.Parallel()
|
||||
|
||||
req, err := http.NewRequestWithContext(
|
||||
context.Background(),
|
||||
http.MethodPost,
|
||||
"https://target.example.com/hook",
|
||||
http.NoBody,
|
||||
)
|
||||
require.NoError(t, err)
|
||||
|
||||
names := delivery.ExportApplyRequestHeaders(
|
||||
req,
|
||||
&database.Event{
|
||||
Headers: `{"X-Custom":["value1"],"X-Broken":"not a list"}`,
|
||||
},
|
||||
&delivery.HTTPTargetConfig{},
|
||||
"webhooker/dev",
|
||||
)
|
||||
|
||||
assert.Empty(t, names)
|
||||
assert.Empty(t, req.Header.Get("X-Custom"))
|
||||
}
|
||||
|
||||
func TestProcessDelivery_RoutesToCorrectHandler(
|
||||
t *testing.T,
|
||||
) {
|
||||
|
||||
@@ -376,3 +376,97 @@ func TestFailedResultWriteLeavesDeliveryRecoverable(
|
||||
database.DeliveryStatusPending,
|
||||
)
|
||||
}
|
||||
|
||||
// TestFailedResultWriteWithRetriesLeavesDeliveryRecoverable is the same
|
||||
// rule for a target with retries: whatever the receiver answered, the
|
||||
// delivery stays pending and no retry is scheduled. The circuit breaker
|
||||
// still learns the answer, because it describes the target's health,
|
||||
// not the database's.
|
||||
func TestFailedResultWriteWithRetriesLeavesDeliveryRecoverable(
|
||||
t *testing.T,
|
||||
) {
|
||||
t.Parallel()
|
||||
|
||||
// The "send succeeded" case starts with the breaker tripped open,
|
||||
// so the delivery goes out as its probe and only a recorded
|
||||
// success closes it again.
|
||||
tests := []struct {
|
||||
name string
|
||||
answer int
|
||||
tripped bool
|
||||
wantBreaker delivery.CircuitState
|
||||
}{
|
||||
{"send succeeded", http.StatusOK, true, delivery.CircuitClosed},
|
||||
{"send failed", http.StatusBadGateway, false, delivery.CircuitOpen},
|
||||
}
|
||||
|
||||
for _, tc := range tests {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
s := newISetup(t)
|
||||
targetID := uuid.New().String()
|
||||
|
||||
ts := httptest.NewServer(http.HandlerFunc(
|
||||
func(w http.ResponseWriter, _ *http.Request) {
|
||||
w.WriteHeader(tc.answer)
|
||||
},
|
||||
))
|
||||
defer ts.Close()
|
||||
|
||||
event := iSeedEvent(
|
||||
t, s.WebhookDB, s.WebhookID, `{"unwritable":true}`,
|
||||
)
|
||||
|
||||
d := iSeedDelivery(
|
||||
t, s.WebhookDB, event.ID, targetID,
|
||||
database.DeliveryStatusPending,
|
||||
)
|
||||
|
||||
require.NoError(
|
||||
t,
|
||||
s.WebhookDB.Exec("drop table delivery_results").Error,
|
||||
)
|
||||
|
||||
// A single failure opens this breaker, and with no
|
||||
// cooldown an open breaker lets the next delivery
|
||||
// through as a probe.
|
||||
cb := delivery.NewTestCircuitBreaker(1, 0)
|
||||
if tc.tripped {
|
||||
cb.RecordFailure()
|
||||
}
|
||||
|
||||
s.Engine.ExportSetCircuitBreaker(targetID, cb)
|
||||
|
||||
full := &database.Delivery{
|
||||
EventID: event.ID,
|
||||
TargetID: targetID,
|
||||
Status: database.DeliveryStatusPending,
|
||||
Event: event,
|
||||
Target: database.Target{
|
||||
Name: "unwritable",
|
||||
Type: database.TargetTypeHTTP,
|
||||
Config: iHTTPConfig(ts.URL),
|
||||
MaxRetries: 3,
|
||||
},
|
||||
}
|
||||
full.ID = d.ID
|
||||
|
||||
sched := &recordingScheduler{}
|
||||
|
||||
s.Engine.ExportDeliverHTTPWithScheduler(
|
||||
context.Background(), s.WebhookDB, full,
|
||||
&delivery.Task{
|
||||
DeliveryID: d.ID,
|
||||
TargetID: targetID,
|
||||
AttemptNum: 1,
|
||||
},
|
||||
sched,
|
||||
)
|
||||
|
||||
iAssertStatus(t, s.WebhookDB, d.ID, database.DeliveryStatusPending)
|
||||
assert.Empty(t, sched.delays, "no retry may be scheduled")
|
||||
assert.Equal(t, tc.wantBreaker, cb.State())
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -179,6 +179,27 @@ func TestDoHTTPRequest_TransportErrorMasksURL(t *testing.T) {
|
||||
)
|
||||
}
|
||||
|
||||
// TestDoHTTPRequest_UnparsableURLIsMasked is the same for an HTTP
|
||||
// target URL that no request can be built from.
|
||||
func TestDoHTTPRequest_UnparsableURLIsMasked(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
e := testEngine(t, 1)
|
||||
|
||||
statusCode, _, _, reqErr := e.ExportDoHTTPRequest(
|
||||
context.TODO(),
|
||||
&delivery.HTTPTargetConfig{
|
||||
URL: "https://hooks.example.com" + maskSecretPath + "\n",
|
||||
},
|
||||
&database.Event{},
|
||||
)
|
||||
require.Error(t, reqErr)
|
||||
assert.Zero(t, statusCode)
|
||||
|
||||
assertNoCredential(t, reqErr.Error())
|
||||
assert.Contains(t, reqErr.Error(), "invalid control character")
|
||||
}
|
||||
|
||||
// TestValidateTargetURL_UnparsableURLIsMasked proves the SSRF
|
||||
// validator's error does not carry the submitted URL, which
|
||||
// the handler both logs and shows.
|
||||
|
||||
@@ -14,18 +14,16 @@ import (
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
// minNonTestFiles guards the walk below against passing because it
|
||||
// found nothing to look at. The tree held 60 non-test .go files when
|
||||
// this was written.
|
||||
const minNonTestFiles = 40
|
||||
|
||||
// isRowProducer reports whether name is a method that returns a
|
||||
// database/sql row handle. GORM's Row and Rows return *sql.Row and
|
||||
// *sql.Rows, so Scan on the result of one of them is database/sql's
|
||||
// Scan and never (*gorm.DB).Scan.
|
||||
// isRowProducer reports whether name is GORM's Row or database/sql's
|
||||
// QueryRow or QueryRowContext, which return a *sql.Row whose Scan is
|
||||
// database/sql's and not (*gorm.DB).Scan. GORM's Rows is not listed:
|
||||
// it also returns an error, so Scan is never called on its result
|
||||
// directly. It matches the method name only and resolves no types, so
|
||||
// a repo-local method with one of these names that returns *gorm.DB
|
||||
// gets past it: Scan on that method's result is not reported.
|
||||
func isRowProducer(name string) bool {
|
||||
switch name {
|
||||
case "Row", "Rows", "QueryRow", "QueryRowContext":
|
||||
case "Row", "QueryRow", "QueryRowContext":
|
||||
return true
|
||||
default:
|
||||
return false
|
||||
@@ -50,9 +48,14 @@ func receiverIsRowHandle(x ast.Expr) bool {
|
||||
}
|
||||
|
||||
// unguardedScans returns the position of every Scan call in file whose
|
||||
// receiver is not a row handle. It fails closed: a receiver it cannot
|
||||
// resolve syntactically — a local variable, a struct field — is
|
||||
// reported rather than assumed safe.
|
||||
// receiver is not a call to a row producer. It fails closed: any other
|
||||
// receiver — a local variable, a struct field, a call to any other
|
||||
// method — is reported rather than assumed safe.
|
||||
//
|
||||
// It sees only calls written x.Scan(...). A method value, f := db.Scan
|
||||
// followed by f(&v), is out of scope: Scan is never the called
|
||||
// expression there, and nobody writes a query that way by accident,
|
||||
// which is the mistake this check exists to catch.
|
||||
func unguardedScans(
|
||||
fset *token.FileSet, file *ast.File,
|
||||
) []token.Position {
|
||||
@@ -111,15 +114,15 @@ func skipDir(name string) bool {
|
||||
}
|
||||
}
|
||||
|
||||
// walkNonTestGo parses every non-test .go file under root and returns
|
||||
// how many it parsed along with every unguarded Scan it found.
|
||||
func walkNonTestGo(t *testing.T, root string) (int, []string) {
|
||||
// walkNonTestGo parses every non-test .go file under root. It returns
|
||||
// the directories, relative to root, it parsed a file in, along with
|
||||
// every unguarded Scan it found.
|
||||
func walkNonTestGo(t *testing.T, root string) (map[string]bool, []string) {
|
||||
t.Helper()
|
||||
|
||||
var (
|
||||
parsed int
|
||||
hits []string
|
||||
)
|
||||
walked := map[string]bool{}
|
||||
|
||||
var hits []string
|
||||
|
||||
fset := token.NewFileSet()
|
||||
|
||||
@@ -147,7 +150,12 @@ func walkNonTestGo(t *testing.T, root string) (int, []string) {
|
||||
return err
|
||||
}
|
||||
|
||||
parsed++
|
||||
dir, err := filepath.Rel(root, filepath.Dir(path))
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
walked[dir] = true
|
||||
|
||||
for _, pos := range unguardedScans(fset, file) {
|
||||
hits = append(hits, relPosition(root, pos))
|
||||
@@ -157,7 +165,7 @@ func walkNonTestGo(t *testing.T, root string) (int, []string) {
|
||||
},
|
||||
))
|
||||
|
||||
return parsed, hits
|
||||
return walked, hits
|
||||
}
|
||||
|
||||
// isNonTestGo reports whether a file name is Go source this check
|
||||
@@ -189,19 +197,39 @@ func relPosition(root string, pos token.Position) string {
|
||||
// logged with its values interpolated. The package comment states the
|
||||
// limit; this fails when someone adds a call site anyway.
|
||||
//
|
||||
// The current tree has one caller, internal/database/database_test.go,
|
||||
// which this check does not govern: it is test-only and its SELECT 1
|
||||
// binds nothing.
|
||||
// Test files are not governed: what a test binds is fixture data.
|
||||
func TestGormScanIsNeverCalledOutsideTests(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
parsed, offenders := walkNonTestGo(t, moduleRoot(t))
|
||||
root := moduleRoot(t)
|
||||
walked, offenders := walkNonTestGo(t, root)
|
||||
|
||||
// The module's packages are static, templates, and every directory
|
||||
// directly under cmd and internal. Each holds non-test code, so one
|
||||
// the walk parsed nothing in was skipped, and a Scan there would
|
||||
// pass unseen.
|
||||
packages := []string{"static", "templates"}
|
||||
|
||||
for _, parent := range []string{"cmd", "internal"} {
|
||||
entries, err := os.ReadDir(filepath.Join(root, parent))
|
||||
require.NoError(t, err)
|
||||
|
||||
for _, entry := range entries {
|
||||
if !entry.IsDir() {
|
||||
continue
|
||||
}
|
||||
|
||||
packages = append(packages, filepath.Join(parent, entry.Name()))
|
||||
}
|
||||
}
|
||||
|
||||
for _, dir := range packages {
|
||||
require.True(
|
||||
t, walked[dir],
|
||||
"the walk parsed no non-test .go file in %s", dir,
|
||||
)
|
||||
}
|
||||
|
||||
require.GreaterOrEqual(
|
||||
t, parsed, minNonTestFiles,
|
||||
"parsed %d non-test .go files, so this check found "+
|
||||
"nothing to look at", parsed,
|
||||
)
|
||||
require.Empty(
|
||||
t, offenders,
|
||||
"Scan called on a receiver this check cannot show is a "+
|
||||
@@ -222,18 +250,51 @@ type scanGuardCase struct {
|
||||
want int
|
||||
}
|
||||
|
||||
// scanGuardCases covers each receiver form unguardedScans names, plus
|
||||
// each row producer isRowProducer lets through. Each body is valid Go
|
||||
// inside plantedFile.
|
||||
func scanGuardCases() []scanGuardCase {
|
||||
return []scanGuardCase{
|
||||
{"gorm chain", `db.DB().Raw("SELECT 1").Scan(&v)`, 1},
|
||||
{"gorm receiver", `gdb.Scan(&v)`, 1},
|
||||
{"gorm via variable", "q := gdb.Raw(\"x\")\nq.Scan(&v)", 1},
|
||||
{"gorm model chain", `gdb.Model(&x).Scan(&v)`, 1},
|
||||
{"sql row", `gdb.Raw("SELECT 1").Row().Scan(&v)`, 0},
|
||||
{"sql rows", `gdb.Raw("SELECT 1").Rows().Scan(&v)`, 0},
|
||||
{"local variable", "q := gdb.Raw(\"SELECT 1\")\n\tq.Scan(&v)", 1},
|
||||
{"struct field", `s.db.Scan(&v)`, 1},
|
||||
{"gorm chain", `gdb.Raw("SELECT 1").Scan(&v)`, 1},
|
||||
{
|
||||
"sql rows in a variable",
|
||||
"rows, _ := gdb.Raw(\"SELECT 1\").Rows()\n\trows.Scan(&v)",
|
||||
1,
|
||||
},
|
||||
{"gorm Row", `gdb.Raw("SELECT 1").Row().Scan(&v)`, 0},
|
||||
{"sql QueryRow", `sqlDB.QueryRow("SELECT 1").Scan(&v)`, 0},
|
||||
{
|
||||
"sql QueryRowContext",
|
||||
`sqlDB.QueryRowContext(ctx, "SELECT 1").Scan(&v)`,
|
||||
0,
|
||||
},
|
||||
{"unrelated call", `gdb.Find(&v)`, 0},
|
||||
}
|
||||
}
|
||||
|
||||
// plantedFile wraps one case body in a function that declares every
|
||||
// name the bodies use, so each body is the Go it stands for. The result
|
||||
// is parsed, never compiled.
|
||||
const plantedFile = `package p
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
type store struct{ db *gorm.DB }
|
||||
|
||||
func f(ctx context.Context, gdb *gorm.DB, sqlDB *sql.DB, s store) {
|
||||
var v int
|
||||
|
||||
%s
|
||||
}
|
||||
`
|
||||
|
||||
// TestScanGuard_ReportsPlantedCalls proves the check fires. Without it
|
||||
// a detector that matched nothing would satisfy the walk above no
|
||||
// matter what the tree contained.
|
||||
@@ -245,9 +306,7 @@ func TestScanGuard_ReportsPlantedCalls(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
fset := token.NewFileSet()
|
||||
src := fmt.Sprintf(
|
||||
"package p\n\nfunc f() {\n\t%s\n}\n", tc.body,
|
||||
)
|
||||
src := fmt.Sprintf(plantedFile, tc.body)
|
||||
|
||||
file, err := parser.ParseFile(
|
||||
fset, tc.name+".go", src, 0,
|
||||
|
||||
@@ -15,10 +15,12 @@ import (
|
||||
// eventBodyQuery reads one event's stored body as bytes. The cast
|
||||
// to blob is what makes the driver hand back the stored bytes
|
||||
// rather than a string conversion, so Content-Length taken from
|
||||
// the result matches what goes on the wire. The soft-delete
|
||||
// predicate is spelled out because Raw bypasses GORM's default
|
||||
// scope, and it is what stops a reaped event still being
|
||||
// downloadable.
|
||||
// the result matches what goes on the wire. The retention reaper
|
||||
// deletes event rows outright, so a reaped event is simply gone
|
||||
// and the query finds no row. The deleted_at predicate repeats
|
||||
// the soft-delete scope GORM adds to its own queries, which Raw
|
||||
// bypasses; nothing soft-deletes an event, so today it excludes
|
||||
// nothing.
|
||||
const eventBodyQuery = "SELECT cast(body as blob) " +
|
||||
"FROM events WHERE id = ? AND webhook_id = ? AND deleted_at IS NULL"
|
||||
|
||||
|
||||
@@ -405,10 +405,11 @@ func TestHandleEventBodyDownload_UnknownEvent404s(t *testing.T) {
|
||||
// route. The body is read in one query before any header is
|
||||
// written, so a reaped event cannot produce a partial download:
|
||||
// it is a clean 404 with no Content-Length and no
|
||||
// Content-Disposition. Both removals the codebase performs are
|
||||
// covered — the reaper hard-deletes, and a soft-deleted row is
|
||||
// excluded by the query's own deleted_at predicate rather than
|
||||
// by GORM's default scope, which Raw bypasses.
|
||||
// Content-Disposition. The reaper deletes event rows outright,
|
||||
// which is the "hard deleted" case. The "soft deleted" case
|
||||
// covers a row no code produces today: it only pins the query's
|
||||
// own deleted_at predicate, the soft-delete condition Raw would
|
||||
// otherwise skip.
|
||||
func TestHandleEventBodyDownload_ReapedEvent404s(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
|
||||
@@ -145,8 +145,9 @@ func (h *Handlers) resubmitEvent(
|
||||
// per-webhook database files — a sibling webhook's event is not in the
|
||||
// database being queried at all — and is there so the scoping survives
|
||||
// any future change that puts more than one webhook's events in one
|
||||
// file. Going through Model applies GORM's soft-delete scope, which is
|
||||
// what stops a reaped event being resubmitted.
|
||||
// file. A reaped event is not found because the retention reaper
|
||||
// deletes its row outright rather than marking it deleted; see
|
||||
// deleteEvents in internal/database/retention.go.
|
||||
func loadResubmitSource(
|
||||
webhookDB *gorm.DB,
|
||||
webhookID, eventID string,
|
||||
|
||||
@@ -272,10 +272,12 @@ func requestEventSource(
|
||||
|
||||
// createAndFanOut writes the event and one pending delivery per target,
|
||||
// and adds them to the webhook's running totals, in a single
|
||||
// transaction, then hands the tasks to the delivery engine. It is the
|
||||
// only path by which an event and its deliveries are created, so a
|
||||
// resubmitted event is retried, SSRF-guarded and circuit-broken
|
||||
// exactly as a received one is.
|
||||
// transaction, then hands the tasks to the delivery engine. Every
|
||||
// event is created here, received or resubmitted, so a resubmitted
|
||||
// event is retried, SSRF-guarded and circuit-broken exactly as a
|
||||
// received one is. Per-delivery replay is the one other path that
|
||||
// creates a delivery: it adds one to an existing event without
|
||||
// coming through here.
|
||||
//
|
||||
// The tasks are returned as well as queued, so a caller can report how
|
||||
// many targets the event went to.
|
||||
|
||||
@@ -1,6 +1,8 @@
|
||||
package server
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"testing"
|
||||
|
||||
@@ -37,6 +39,14 @@ func SentryClientOptionsForTest(
|
||||
return sentryClientOptions(dsn, release)
|
||||
}
|
||||
|
||||
// CleanShutdownForTest runs the server's stop hook, cleanShutdown,
|
||||
// against hs: a server the test started itself, so it can hold a
|
||||
// request open across the drain. Sentry is off.
|
||||
func CleanShutdownForTest(ctx context.Context, hs *http.Server) {
|
||||
s := &Server{log: slog.New(slog.DiscardHandler), httpServer: hs}
|
||||
s.cleanShutdown(ctx)
|
||||
}
|
||||
|
||||
// newServerForTest builds a Server through New, as the application
|
||||
// does, on a lifecycle that is never started: the hooks New adds to
|
||||
// it never run, so nothing listens.
|
||||
|
||||
@@ -0,0 +1,215 @@
|
||||
package server_test
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"net/url"
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
// maxResubmits bounds the requests the tests below send to the
|
||||
// resubmit route. The route's rate limit belongs to the middleware;
|
||||
// this only has to sit well above it, so that a route without the
|
||||
// limiter fails its test instead of looping.
|
||||
const maxResubmits = 100
|
||||
|
||||
// resubmitPath is the resubmit route for one stored event.
|
||||
func resubmitPath(webhookID, eventID string) string {
|
||||
return "/hook/" + webhookID + "/events/" + eventID + "/resubmit"
|
||||
}
|
||||
|
||||
// csrfForm is a resubmit form carrying the given CSRF token.
|
||||
func csrfForm(token string) url.Values {
|
||||
form := url.Values{}
|
||||
form.Set("csrf_token", token)
|
||||
|
||||
return form
|
||||
}
|
||||
|
||||
// TestEventResubmit_SignedOutRequestsNeverReachTheRateLimit pins
|
||||
// RequireAuth on the resubmit route. The handler also turns away a
|
||||
// request without a session, with the same redirect, so a refusal
|
||||
// alone would pass without RequireAuth. What RequireAuth adds is that
|
||||
// it refuses such a request before the route's rate limit, so a
|
||||
// signed-out client cannot spend the budget a signed-in user
|
||||
// resubmits from. Each request carries a CSRF token valid for its own
|
||||
// cookie, so CSRF lets it through to RequireAuth.
|
||||
func TestEventResubmit_SignedOutRequestsNeverReachTheRateLimit(
|
||||
t *testing.T,
|
||||
) {
|
||||
t.Parallel()
|
||||
|
||||
env := newTestEnv(t)
|
||||
|
||||
userID, _ := env.seedUser(t, "resubmitter", "somepassword")
|
||||
wh := env.seedWebhook(t, userID)
|
||||
evt := env.seedEvent(t, wh.ID, `{"resubmit":"me"}`)
|
||||
path := resubmitPath(wh.ID, evt.ID)
|
||||
logsPath := "/hook/" + wh.ID + "/events"
|
||||
|
||||
token, signedOut := env.csrfFrom(t, "/pages/login", nil)
|
||||
|
||||
for i := range maxResubmits {
|
||||
w := env.post(path, csrfForm(token), signedOut)
|
||||
require.Equal(t, http.StatusSeeOther, w.Code, "request %d", i)
|
||||
require.Equal(
|
||||
t, "/pages/login", w.Header().Get("Location"),
|
||||
"request %d", i,
|
||||
)
|
||||
}
|
||||
|
||||
require.Equal(
|
||||
t, int64(1), env.countEvents(t, wh.ID),
|
||||
"a signed-out request must store nothing",
|
||||
)
|
||||
|
||||
token, cookies := env.csrfFrom(
|
||||
t, logsPath, env.authCookies(t, userID, "resubmitter"),
|
||||
)
|
||||
|
||||
env.requireNotice(
|
||||
t, env.post(path, csrfForm(token), cookies),
|
||||
logsPath, "resubmit-no-targets",
|
||||
"this source has no active targets", cookies,
|
||||
)
|
||||
}
|
||||
|
||||
// TestEventResubmit_RefusedWithoutAValidCSRFToken pins CSRF on the
|
||||
// resubmit route: a signed-in user's POST is refused with 403, and
|
||||
// stores nothing, unless it carries the token issued to that user's
|
||||
// own browser.
|
||||
func TestEventResubmit_RefusedWithoutAValidCSRFToken(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
env := newTestEnv(t)
|
||||
|
||||
userID, _ := env.seedUser(t, "resubmitter", "somepassword")
|
||||
wh := env.seedWebhook(t, userID)
|
||||
evt := env.seedEvent(t, wh.ID, `{"resubmit":"me"}`)
|
||||
path := resubmitPath(wh.ID, evt.ID)
|
||||
logsPath := "/hook/" + wh.ID + "/events"
|
||||
|
||||
token, cookies := env.csrfFrom(
|
||||
t, logsPath, env.authCookies(t, userID, "resubmitter"),
|
||||
)
|
||||
otherBrowsers, _ := env.csrfFrom(t, "/pages/login", nil)
|
||||
|
||||
for name, form := range map[string]url.Values{
|
||||
"no token": {},
|
||||
"a malformed token": csrfForm("not-a-token"),
|
||||
"another browser's token": csrfForm(otherBrowsers),
|
||||
} {
|
||||
assert.Equal(
|
||||
t, http.StatusForbidden,
|
||||
env.post(path, form, cookies).Code, name,
|
||||
)
|
||||
}
|
||||
|
||||
assert.Equal(
|
||||
t, int64(1), env.countEvents(t, wh.ID),
|
||||
"a refused request must store nothing",
|
||||
)
|
||||
|
||||
// The same request with the user's own token goes through, so the
|
||||
// refusals above were the token's doing.
|
||||
env.requireNotice(
|
||||
t, env.post(path, csrfForm(token), cookies),
|
||||
logsPath, "resubmit-no-targets",
|
||||
"this source has no active targets", cookies,
|
||||
)
|
||||
}
|
||||
|
||||
// TestEventResubmit_AnotherWebhooksEvent404s pins, on the route as
|
||||
// registered, that a signed-in user gets 404, and nothing is stored,
|
||||
// for an event of a webhook another user owns, which the handler's
|
||||
// ownership check refuses, and for another webhook's event posted
|
||||
// under a webhook the user does own, which the event lookup refuses.
|
||||
// The user's own event, posted the same way, is accepted, so the
|
||||
// second 404 comes from the lookup and not from a route that never
|
||||
// passed the event ID to the handler.
|
||||
func TestEventResubmit_AnotherWebhooksEvent404s(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
env := newTestEnv(t)
|
||||
|
||||
ownerID, _ := env.seedUser(t, "owner", "somepassword")
|
||||
owners := env.seedWebhook(t, ownerID)
|
||||
ownersEvent := env.seedEvent(t, owners.ID, `{"owner":"only"}`)
|
||||
|
||||
intruderID, _ := env.seedUser(t, "intruder", "somepassword")
|
||||
intruders := env.seedWebhook(t, intruderID)
|
||||
intrudersEvent := env.seedEvent(
|
||||
t, intruders.ID, `{"intruder":"own"}`,
|
||||
)
|
||||
intrudersLogs := "/hook/" + intruders.ID + "/events"
|
||||
|
||||
token, cookies := env.csrfFrom(
|
||||
t, intrudersLogs, env.authCookies(t, intruderID, "intruder"),
|
||||
)
|
||||
|
||||
for name, path := range map[string]string{
|
||||
"another user's webhook": resubmitPath(
|
||||
owners.ID, ownersEvent.ID,
|
||||
),
|
||||
"another webhook's event": resubmitPath(
|
||||
intruders.ID, ownersEvent.ID,
|
||||
),
|
||||
} {
|
||||
w := env.post(path, csrfForm(token), cookies)
|
||||
assert.Equal(t, http.StatusNotFound, w.Code, name)
|
||||
}
|
||||
|
||||
assert.Equal(t, int64(1), env.countEvents(t, owners.ID))
|
||||
assert.Equal(t, int64(1), env.countEvents(t, intruders.ID))
|
||||
|
||||
env.requireNotice(
|
||||
t,
|
||||
env.post(
|
||||
resubmitPath(intruders.ID, intrudersEvent.ID),
|
||||
csrfForm(token), cookies,
|
||||
),
|
||||
intrudersLogs, "resubmit-no-targets",
|
||||
"this source has no active targets", cookies,
|
||||
)
|
||||
}
|
||||
|
||||
// TestEventResubmit_RateLimited pins the rate limit on the resubmit
|
||||
// route: a signed-in user's resubmits are accepted until the budget
|
||||
// is spent, and then refused with 429.
|
||||
func TestEventResubmit_RateLimited(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
env := newTestEnv(t)
|
||||
|
||||
userID, _ := env.seedUser(t, "resubmitter", "somepassword")
|
||||
wh := env.seedWebhook(t, userID)
|
||||
evt := env.seedEvent(t, wh.ID, `{"resubmit":"me"}`)
|
||||
path := resubmitPath(wh.ID, evt.ID)
|
||||
|
||||
token, cookies := env.csrfFrom(
|
||||
t, "/hook/"+wh.ID+"/events",
|
||||
env.authCookies(t, userID, "resubmitter"),
|
||||
)
|
||||
|
||||
limited := false
|
||||
|
||||
for range maxResubmits {
|
||||
code := env.post(path, csrfForm(token), cookies).Code
|
||||
if code == http.StatusTooManyRequests {
|
||||
limited = true
|
||||
|
||||
break
|
||||
}
|
||||
|
||||
require.Equal(
|
||||
t, http.StatusSeeOther, code,
|
||||
"a resubmit within the budget must be accepted",
|
||||
)
|
||||
}
|
||||
|
||||
assert.True(
|
||||
t, limited, "repeated resubmits must eventually be refused",
|
||||
)
|
||||
}
|
||||
@@ -444,6 +444,23 @@ func (e *testEnv) countDeliveries(
|
||||
return count
|
||||
}
|
||||
|
||||
// countEvents reports how many events a webhook's database holds.
|
||||
func (e *testEnv) countEvents(t *testing.T, webhookID string) int64 {
|
||||
t.Helper()
|
||||
|
||||
webhookDB, err := e.dbMgr.GetDB(webhookID)
|
||||
require.NoError(t, err)
|
||||
|
||||
var count int64
|
||||
|
||||
require.NoError(
|
||||
t,
|
||||
webhookDB.Model(&database.Event{}).Count(&count).Error,
|
||||
)
|
||||
|
||||
return count
|
||||
}
|
||||
|
||||
// storedHash reads the current password hash for a username.
|
||||
func (e *testEnv) storedHash(t *testing.T, username string) string {
|
||||
t.Helper()
|
||||
|
||||
@@ -39,6 +39,12 @@ const (
|
||||
// refuses to spend, leaving it for the hooks that run after the
|
||||
// server: the delivery engine, the healthcheck, the webhook DB
|
||||
// manager and the database close.
|
||||
//
|
||||
// Its value is not tuned to those hooks, which take about a
|
||||
// millisecond between them. It is what the 5s fx stop timeout in
|
||||
// cmd/webhooker leaves after a full ShutdownTimeout drain, so a
|
||||
// drain that starts on a full budget still gets all of
|
||||
// ShutdownTimeout.
|
||||
TailHookReserve = 2 * time.Second
|
||||
|
||||
// sentryFlushTimeout is the longest wait for Sentry to flush
|
||||
@@ -59,6 +65,16 @@ const (
|
||||
// key off it, and a zero exit would read as a deliberate stop.
|
||||
const StartupFailureExitCode = 1
|
||||
|
||||
// DrainBudget reports how long the HTTP drain may wait for in-flight
|
||||
// requests when remaining is the time left on the fx stop context as
|
||||
// the server's stop hook starts. The hooks before the server can
|
||||
// already have spent part of the budget, so the drain takes its time
|
||||
// out of what they left, never out of TailHookReserve. Zero or less
|
||||
// means no wait at all.
|
||||
func DrainBudget(remaining time.Duration) time.Duration {
|
||||
return min(ShutdownTimeout, remaining-TailHookReserve)
|
||||
}
|
||||
|
||||
// SentryFlushBudget reports how long the Sentry flush may run when
|
||||
// remaining is the time left on the fx stop context after the HTTP
|
||||
// drain. sentry.Flush takes a bare duration and honours no context,
|
||||
@@ -261,10 +277,17 @@ func (s *Server) cleanupForExit() {
|
||||
s.log.Info("cleaning up")
|
||||
}
|
||||
|
||||
// cleanShutdown drains the HTTP server and flushes Sentry inside what
|
||||
// is left of the fx stop budget. A context carrying no deadline — a
|
||||
// caller outside the fx lifecycle — gets the full ShutdownTimeout.
|
||||
func (s *Server) cleanShutdown(ctx context.Context) {
|
||||
ctxShutdown, shutdownCancel := context.WithTimeout(
|
||||
ctx, ShutdownTimeout,
|
||||
)
|
||||
drain := ShutdownTimeout
|
||||
|
||||
if deadline, ok := ctx.Deadline(); ok {
|
||||
drain = DrainBudget(time.Until(deadline))
|
||||
}
|
||||
|
||||
ctxShutdown, shutdownCancel := context.WithTimeout(ctx, drain)
|
||||
defer shutdownCancel()
|
||||
|
||||
err := s.httpServer.Shutdown(ctxShutdown)
|
||||
|
||||
@@ -1,13 +1,148 @@
|
||||
package server_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net"
|
||||
"net/http"
|
||||
"testing"
|
||||
"testing/synctest"
|
||||
"time"
|
||||
|
||||
"github.com/stretchr/testify/require"
|
||||
"sneak.berlin/go/webhooker/internal/server"
|
||||
)
|
||||
|
||||
// TestDrainBudget covers the clamp that keeps the HTTP drain from
|
||||
// spending the tail hooks' share of the fx stop budget when the hooks
|
||||
// before the server have already used part of it.
|
||||
func TestDrainBudget(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
tests := []struct {
|
||||
name string
|
||||
remaining time.Duration
|
||||
want time.Duration
|
||||
}{
|
||||
{
|
||||
name: "only the reserve is left",
|
||||
remaining: server.TailHookReserve,
|
||||
want: 0,
|
||||
},
|
||||
{
|
||||
name: "earlier hooks spent part of the budget",
|
||||
remaining: server.TailHookReserve + time.Second,
|
||||
want: time.Second,
|
||||
},
|
||||
{
|
||||
name: "capped at the nominal timeout",
|
||||
remaining: time.Hour,
|
||||
want: server.ShutdownTimeout,
|
||||
},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
require.Equal(t, tt.want, server.DrainBudget(tt.remaining))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// TestCleanShutdown_LeavesTailHookReserve stops the server with a
|
||||
// request still in flight, after the hooks before it have spent all
|
||||
// of the stop budget but TailHookReserve. The drain must give up at
|
||||
// once rather than wait for the request: what is left belongs to the
|
||||
// hooks after the server, the database close among them. A drain
|
||||
// bounded only by ShutdownTimeout waits until the stop context
|
||||
// expires, and fx then skips those hooks.
|
||||
//
|
||||
// The test runs in a synctest bubble, whose clock moves only while
|
||||
// every goroutine in it is blocked, so a drain that gives up at once
|
||||
// leaves the stop context unexpired however slow the host is. The
|
||||
// request travels over net.Pipe because a goroutine waiting on a
|
||||
// real socket would stop that clock from moving at all.
|
||||
func TestCleanShutdown_LeavesTailHookReserve(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
synctest.Test(t, func(t *testing.T) {
|
||||
entered := make(chan struct{})
|
||||
release := make(chan struct{})
|
||||
|
||||
hs := &http.Server{
|
||||
Handler: http.HandlerFunc(
|
||||
func(http.ResponseWriter, *http.Request) {
|
||||
close(entered)
|
||||
<-release
|
||||
},
|
||||
),
|
||||
ReadHeaderTimeout: time.Second,
|
||||
}
|
||||
|
||||
srvConn, cliConn := net.Pipe()
|
||||
|
||||
listener := pipeListener{
|
||||
conns: make(chan net.Conn, 1),
|
||||
closed: make(chan struct{}),
|
||||
}
|
||||
listener.conns <- srvConn
|
||||
|
||||
go func() { _ = hs.Serve(listener) }()
|
||||
|
||||
// Cleanups run last first: the handler returns, then closing
|
||||
// the client end ends the server's write of the response.
|
||||
t.Cleanup(func() { _ = cliConn.Close() })
|
||||
t.Cleanup(func() { close(release) })
|
||||
|
||||
_, err := cliConn.Write(
|
||||
[]byte("GET / HTTP/1.1\r\nHost: webhooker.test\r\n\r\n"),
|
||||
)
|
||||
require.NoError(t, err)
|
||||
|
||||
<-entered
|
||||
|
||||
stopCtx, cancel := context.WithTimeout(
|
||||
t.Context(), server.TailHookReserve,
|
||||
)
|
||||
defer cancel()
|
||||
|
||||
server.CleanShutdownForTest(stopCtx, hs)
|
||||
|
||||
require.NoError(
|
||||
t, stopCtx.Err(), "the drain spent the tail hooks' reserve",
|
||||
)
|
||||
})
|
||||
}
|
||||
|
||||
// pipeListener is the net.Listener http.Server.Serve needs to serve
|
||||
// the server end of a net.Pipe: Accept returns that one connection,
|
||||
// then waits until Close, as a real listener with no more clients
|
||||
// does.
|
||||
type pipeListener struct {
|
||||
conns chan net.Conn
|
||||
closed chan struct{}
|
||||
}
|
||||
|
||||
func (l pipeListener) Accept() (net.Conn, error) {
|
||||
select {
|
||||
case conn := <-l.conns:
|
||||
return conn, nil
|
||||
case <-l.closed:
|
||||
return nil, net.ErrClosed
|
||||
}
|
||||
}
|
||||
|
||||
func (l pipeListener) Close() error {
|
||||
close(l.closed)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// Addr is never called by http.Server.Serve.
|
||||
func (pipeListener) Addr() net.Addr {
|
||||
return nil
|
||||
}
|
||||
|
||||
// TestSentryFlushBudget covers the clamp that keeps the Sentry flush
|
||||
// from spending the tail hooks' share of the fx stop budget.
|
||||
// sentry.Flush ignores the stop context, so without the clamp a
|
||||
|
||||
Reference in New Issue
Block a user