Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,12 @@

### Added

- `GET /healthz` on the webhook source's address. It answers 200 while the
source admits deliveries and 503 once it is closing, needs no signature,
and is not counted in `webhook_requests_total`. A platform that routes one
port to a service checks health on that port, and the pipeline's `/healthz`
is on the metrics listener, so a webhook pipeline on Render had no health
check at all: it was marked live because its port was open.
- `sqlflow rollup` measures: `numeric: double` on a `sum`, and a `last` type.
Every sum was cast to `bigint`, so a fractional value was rounded to an
integer at every grain. `last` keeps the value of the bucket's latest finer
Expand Down
9 changes: 9 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -820,6 +820,15 @@ body bound applies before the body is read and before the signature is
checked, so an unsigned oversized request never holds memory past it. The
default is 25 MiB, GitHub's payload ceiling.

`GET /healthz` on the same address answers `{"status":"ok"}` while the source
admits deliveries, and 503 once it is closing. It needs no signature and reads
no body, so a platform's health check or an uptime monitor can call it, and
`HEAD` works too. A queue that is full is backpressure, not failure, and still
answers 200. Health checks are not counted in `webhook_requests_total`, which
counts deliveries. The pipeline's own `/healthz`, on the metrics listener, is
unchanged: this one exists because a platform that routes a single port to a
service, as Render does, checks health on that port.

The listener holds at most `max_connections` open connections; past that, a
new connection waits in the kernel backlog until one closes. Three fixed
timeouts release a slot: 10s for a connection to send its request headers,
Expand Down
30 changes: 30 additions & 0 deletions internal/webhook/metrics_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -151,3 +151,33 @@ func TestSinkRetry_MetricsNoProviderRecordsNothing(t *testing.T) {
defer resp.Body.Close()
assert.Equal(t, http.StatusOK, resp.StatusCode)
}

// webhook_requests_total is how an operator counts deliveries. A platform
// checks health every few seconds, which would bury them under 200s that
// delivered nothing.
func TestSinkRetry_MetricsDoNotCountHealthChecks(t *testing.T) {
coverage.Covers(t, "sink.retry")
reader := sdkmetric.NewManualReader()
s, err := NewSource(WithMeterProvider(sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader))))
assert.NoError(t, err)
defer s.Close()
drain(s)

srv := httptest.NewServer(s.Handler())
defer srv.Close()

for range 5 {
resp, err := http.Get(srv.URL + "/healthz")
assert.NoError(t, err)
resp.Body.Close()
}
resp := post(t, srv.URL+"/events", []byte(`{"a":1}`), "", "")
resp.Body.Close()

sum := collect(t, reader, "webhook_requests_total").Data.(metricdata.Sum[int64])
var total int64
for _, dp := range sum.DataPoints {
total += dp.Value
}
assert.Equal(t, int64(1), total)
}
38 changes: 37 additions & 1 deletion internal/webhook/source.go
Original file line number Diff line number Diff line change
Expand Up @@ -192,9 +192,45 @@ func (s *Source) MaxConnections() int {
func (s *Source) Handler() http.Handler {
mux := http.NewServeMux()
mux.HandleFunc("POST /events", s.receiveEvents)

root := http.NewServeMux()
// Outside the metrics middleware. webhook_requests_total is how an
// operator counts deliveries, and a platform checks health every few
// seconds: counted, the checks would bury them under 200s that delivered
// nothing. Registered for every method, because beside the catch-all
// below a GET-only pattern would hand POST /healthz to the delivery mux,
// which answers 404 for a path that exists.
root.HandleFunc("/healthz", s.healthz)
// Wrapped rather than applied per route, so unrouted requests are counted
// as they are by the Python middleware.
return s.metrics.middleware(mux)
root.Handle("/", s.metrics.middleware(mux))
return root
}

// healthz says whether a delivery sent now would be admitted. It is on the
// webhook's own listener because a platform routes one port to a service and
// checks health on that port: the pipeline's /healthz is on the metrics
// listener, which a platform that exposes only this one cannot reach.
//
// It reads no body and checks no signature, since a health check has neither,
// and it admits nothing to the pipeline. It does not wait on the queue: a full
// queue is backpressure, and a platform that reads busy as dead restarts an
// instance while it holds a sender's event. A closing source answers 503, as
// it does to a delivery, so the platform stops routing to it.
func (s *Source) healthz(w http.ResponseWriter, r *http.Request) {
// Monitors send HEAD, and the status line is the whole answer. net/http
// drops the body of a HEAD response.
if r.Method != http.MethodGet && r.Method != http.MethodHead {
w.Header().Set("Allow", "GET, HEAD")
writeJSON(w, http.StatusMethodNotAllowed, `{"detail":"Method Not Allowed"}`)
return
}
select {
case <-s.done:
writeJSON(w, http.StatusServiceUnavailable, `{"detail":"Source is closed"}`)
default:
writeJSON(w, http.StatusOK, `{"status":"ok"}`)
}
}

// Addr reports the bound address, which only differs from the configured one
Expand Down
87 changes: 86 additions & 1 deletion internal/webhook/source_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -225,7 +225,9 @@ func TestSourceWebhook_CloseReleasesBlockedRequest(t *testing.T) {
}
}

func TestSourceWebhook_RoutesOnlyPostEvents(t *testing.T) {
// Deliveries are POST /events and nothing else. /healthz is the one other
// route, and has its own tests.
func TestSourceWebhook_RefusesOtherMethodsAndPaths(t *testing.T) {
coverage.Covers(t, "source.webhook")
s, err := NewSource()
assert.NoError(t, err)
Expand Down Expand Up @@ -506,3 +508,86 @@ func TestSourceWebhook_DefaultsMaxConnectionsTo64(t *testing.T) {
defer s.Close()
assert.Equal(t, 64, s.MaxConnections())
}

// A platform's health check and an uptime monitor have no secret, and must not
// need one: the route reads no body and admits nothing to the pipeline.
func TestSourceWebhook_HealthzAnswersWithoutASignature(t *testing.T) {
coverage.Covers(t, "source.webhook")
s, err := NewSource(WithHMAC(&HMAC{Header: "X-HMAC-Signature", SigKey: "sha256", Secret: "test_secret"}))
assert.NoError(t, err)
defer s.Close()

srv := httptest.NewServer(s.Handler())
defer srv.Close()

resp, err := http.Get(srv.URL + "/healthz")
assert.NoError(t, err)
assert.Equal(t, http.StatusOK, resp.StatusCode)
assert.Equal(t, "application/json", resp.Header.Get("Content-Type"))
assert.Equal(t, `{"status":"ok"}`, readBody(t, resp))

// Monitors send HEAD, and the status line is the whole answer.
resp, err = http.Head(srv.URL + "/healthz")
assert.NoError(t, err)
resp.Body.Close()
assert.Equal(t, http.StatusOK, resp.StatusCode)

resp = post(t, srv.URL+"/healthz", []byte("{}"), "", "")
resp.Body.Close()
assert.Equal(t, http.StatusMethodNotAllowed, resp.StatusCode)

// The route opens nothing else: an unsigned delivery is still refused.
resp = post(t, srv.URL+"/events", []byte("{}"), "", "")
resp.Body.Close()
assert.Equal(t, http.StatusBadRequest, resp.StatusCode)
}

// Healthy means deliveries are being admitted. A source that is closing
// answers 503 to a delivery, so it says the same to the check, and the
// platform stops routing to an instance that would refuse what it was sent.
func TestSourceWebhook_HealthzReportsAClosingSource(t *testing.T) {
coverage.Covers(t, "source.webhook")
s, err := NewSource()
assert.NoError(t, err)

srv := httptest.NewServer(s.Handler())
defer srv.Close()

assert.NoError(t, s.Close())

resp, err := http.Get(srv.URL + "/healthz")
assert.NoError(t, err)
assert.Equal(t, http.StatusServiceUnavailable, resp.StatusCode)
assert.Equal(t, `{"detail":"Source is closed"}`, readBody(t, resp))
}

// A queue that is full is backpressure, not failure. The check must answer
// while a delivery waits on the pipeline, or a platform restarts an instance
// for being busy, at the moment it is holding a sender's event.
func TestSourceWebhook_HealthzAnswersWhileADeliveryWaits(t *testing.T) {
coverage.Covers(t, "source.webhook")
s, err := NewSource()
assert.NoError(t, err)

srv := httptest.NewServer(s.Handler())
// The source closes first, which releases the delivery left waiting.
// The server's Close waits for that delivery and would never return.
defer srv.Close()
defer s.Close()

resp := post(t, srv.URL+"/events", []byte("first"), "", "")
resp.Body.Close()
go func() {
if resp, err := http.Post(srv.URL+"/events", "application/json", strings.NewReader("second")); err == nil {
resp.Body.Close()
}
}()
// Give the second delivery time to reach the full queue.
time.Sleep(100 * time.Millisecond)

client := &http.Client{Timeout: 2 * time.Second}
resp, err = client.Get(srv.URL + "/healthz")
assert.NoError(t, err)
resp.Body.Close()
assert.Equal(t, http.StatusOK, resp.StatusCode)
}
Loading