Skip to content
24 changes: 24 additions & 0 deletions internal/teams/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,10 @@ const (
// inbound Activities cannot spawn unbounded work. The handler acquires a slot
// before acking and sheds load with 503 (the platform retries) when full.
maxConcurrentDispatch = 256
// maxConcurrentReads bounds concurrent inbound request processing (the body
// read + JSON unmarshal before JWT verification), capping peak pre-auth read
// memory at maxConcurrentReads times maxRequestBytes. Mirrors github/gitlab.
maxConcurrentReads = 16
)

// conversation holds what a reply needs: the serviceUrl to POST to and the bot's
Expand Down Expand Up @@ -103,17 +107,37 @@ func serve(srv *http.Server, ln net.Listener, done func(error)) {
}

func (a *adapter) handleMessages(dispatchCtx context.Context, w http.ResponseWriter, r *http.Request, deps core.AdapterDeps) {
// Bound concurrent inbound reads before buffering + unmarshaling the body: a
// flood of large bodies is shed with 503 (the platform retries) rather than
// buffering unbounded memory. readSem is allocated once in newAdapter, so this
// direct read races nothing. Separate from dispatchSem, which bounds post-ack
// dispatch.
select {
case a.readSem <- struct{}{}:
default:
a.log().Warn("teams: inbound concurrency limit reached; shedding with 503", "limit", maxConcurrentReads)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[low] Teams holds readSem through JWT/JWKS validation, unlike the HMAC-fast siblings it copies

The readSem slot is acquired at the top of handleMessages and released only when the function returns (defer func() { <-a.readSem }()), so it is held across validateInbound — RS256 verification plus a potential cold JWKS fetch bounded by jwksFetchTimeout (~15s). In the GitHub/GitLab siblings this constant (maxConcurrentReads = 16) gates only a fast HMAC/token compare, so a slot frees almost immediately. In Teams, under a legitimate burst that coincides with a JWKS refresh, up to 16 slots can stay occupied for seconds, shedding the 17th+ authentic request with 503. The behavior is safe (the platform retries) but the reused 16 doesn't account for the much longer hold, so this is a genuine divergence from the scaffolding it claims to mirror. Consider either releasing readSem right after the body read/unmarshal (before validateInbound) or documenting/sizing the constant for the JWT path. Marking low/uncertain — it only bites during a cold JWKS fetch under concurrent load.

🤖 AI prompt to fix (review before running)
Fix an issue found during code review of lao/botbooter (PR #48).

File: internal/teams/server.go:118  (side RIGHT)
Severity: low
Issue: Teams holds readSem through JWT/JWKS validation, unlike the HMAC-fast siblings it copies

The `readSem` slot is acquired at the top of `handleMessages` and released only when the function returns (`defer func() { <-a.readSem }()`), so it is held across `validateInbound` — RS256 verification plus a potential *cold* JWKS fetch bounded by `jwksFetchTimeout` (~15s). In the GitHub/GitLab siblings this constant (`maxConcurrentReads = 16`) gates only a fast HMAC/token compare, so a slot frees almost immediately. In Teams, under a legitimate burst that coincides with a JWKS refresh, up to 16 slots can stay occupied for seconds, shedding the 17th+ authentic request with 503. The behavior is safe (the platform retries) but the reused `16` doesn't account for the much longer hold, so this is a genuine divergence from the scaffolding it claims to mirror. Consider either releasing `readSem` right after the body read/unmarshal (before `validateInbound`) or documenting/sizing the constant for the JWT path. Marking low/uncertain — it only bites during a cold JWKS fetch under concurrent load.

Apply a minimal, correct fix that resolves this issue. Match the surrounding code's existing style and conventions, and do not change unrelated behavior.

w.WriteHeader(http.StatusServiceUnavailable)
return
}

body, err := io.ReadAll(http.MaxBytesReader(w, r.Body, maxRequestBytes))
if err != nil {
<-a.readSem
w.WriteHeader(http.StatusBadRequest)
return
}

var act inboundActivity
if err := json.Unmarshal(body, &act); err != nil {
<-a.readSem
w.WriteHeader(http.StatusBadRequest)
return
}
// Release the read slot as soon as the body is buffered + unmarshaled: unlike
// the GitHub/GitLab siblings, validateInbound below can block on a cold JWKS
// fetch (up to jwksFetchTimeout), so holding readSem across it would gate
// throughput far more aggressively than the fast HMAC path readSem is sized for.
<-a.readSem

// Authenticate before trusting the body. Use the request context, not runCtx:
// core cancels runCtx before srv.Shutdown drains this handler, so a JWKS
Expand Down
27 changes: 27 additions & 0 deletions internal/teams/server_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import (
"context"
"encoding/json"
"errors"
"io"
"log/slog"
"net"
"net/http"
Expand Down Expand Up @@ -269,6 +270,32 @@ func TestHandleMessages_SaturatedDispatchReturns503(t *testing.T) {
asserts.Equal(t, w2.Code, http.StatusServiceUnavailable, "saturated dispatch sheds load with 503")
}

// readSpy counts Read calls so a test can prove a shed request's body was
// never touched.
type readSpy struct{ reads int }

func (s *readSpy) Read([]byte) (int, error) { s.reads++; return 0, io.EOF }

// A saturated read semaphore must answer 503 before the body is buffered and
// unmarshaled, so a flood of large POSTs cannot exhaust memory ahead of the JWT
// check. Mirrors the github/gitlab sibling test.
func TestHandleMessages_ReadSaturationReturns503(t *testing.T) {
a := testAdapter(t)
a.readSem = make(chan struct{}, 1)
a.readSem <- struct{}{} // occupy the only inbound slot

var got []*core.Message
spy := &readSpy{}
r := httptest.NewRequest(http.MethodPost, a.cfg.Path, spy)
r.Header.Set("Authorization", "Bearer "+mintToken(t, testKID, validClaims(a.cfg.AppID, allowedServiceURL)))
w := httptest.NewRecorder()
a.handleMessages(context.Background(), w, r, captureDeps(&got, nil))

asserts.Equal(t, w.Code, http.StatusServiceUnavailable, "saturated read gate returns 503")
asserts.Equal(t, len(got), 0, "nothing dispatched")
asserts.Equal(t, spy.reads, 0, "shed request's body is never read")
}

func TestHandleMessages_IgnoresNonMessage(t *testing.T) {
a := testAdapter(t)
var got []*core.Message
Expand Down
14 changes: 10 additions & 4 deletions internal/teams/teams.go
Original file line number Diff line number Diff line change
Expand Up @@ -88,10 +88,15 @@ type adapter struct {
// captured at acquire, so a slot a context-ignoring handler never releases is
// confined to that connection rather than leaking for the adapter's lifetime.
dispatchSem chan struct{}
convs map[string]conversation // conversationID -> reply routing info
convOrder []string // FIFO insertion order for bounded eviction
token cachedToken
keys map[string]*jwksKey // kid -> signing key + channel endorsements
// readSem bounds concurrent inbound request processing so a flood of large
// bodies can't buffer unbounded memory before JWT verification (see
// maxConcurrentReads); separate from dispatchSem, which bounds dispatch
// goroutines. Allocated once in newAdapter and never reassigned.
readSem chan struct{}
convs map[string]conversation // conversationID -> reply routing info
convOrder []string // FIFO insertion order for bounded eviction
token cachedToken
keys map[string]*jwksKey // kid -> signing key + channel endorsements
// keysAt is the last JWKS fetch attempt (rate-limits refreshes); keysFreshAt
// is the last successful refresh (drives the jwksMaxAge staleness gate).
keysAt time.Time
Expand Down Expand Up @@ -157,5 +162,6 @@ func newAdapter(cfg Config) (*adapter, error) {
tokenURL: "https://login.microsoftonline.com/" + tenant + "/oauth2/v2.0/token",
openIDURL: openIDConfigURL,
dispatchSem: make(chan struct{}, maxConcurrentDispatch),
readSem: make(chan struct{}, maxConcurrentReads),
}, nil
}
32 changes: 29 additions & 3 deletions internal/whatsapp/cloud/cloud.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ import (
"log/slog"
"net"
"net/http"
"net/url"
"strconv"
"strings"
"sync"
Expand All @@ -52,6 +53,11 @@ const (
// goroutines.
maxConcurrentDispatch = 256

// maxConcurrentReads bounds concurrent inbound request processing (the body
// read before signature verification), capping peak pre-auth read memory at
// maxConcurrentReads times maxRequestBytes. Mirrors the github/gitlab siblings.
maxConcurrentReads = 16

// maxErrorBodyBytes caps how much of a non-2xx response body is read into errors.
maxErrorBodyBytes = 4 << 10 // 4 KiB

Expand Down Expand Up @@ -148,6 +154,11 @@ type adapter struct {
// handler never releases is confined to that connection rather than leaking for
// the adapter's lifetime.
dispatchSem chan struct{}
// readSem bounds concurrent inbound request processing so a flood of large
// bodies can't buffer unbounded memory before signature verification (see
// maxConcurrentReads); separate from dispatchSem, which bounds dispatch
// goroutines. Allocated once in newAdapter and never reassigned.
readSem chan struct{}
}

// log returns the Bot's logger handed over at Connect, or slog.Default()
Expand Down Expand Up @@ -206,7 +217,7 @@ func newAdapter(cfg Config) (*adapter, error) {
if cfg.HTTPClient == nil {
cfg.HTTPClient = &http.Client{Timeout: 30 * time.Second}
}
return &adapter{cfg: cfg, baseURL: graphBaseURL, http: cfg.HTTPClient, dispatchSem: make(chan struct{}, maxConcurrentDispatch)}, nil
return &adapter{cfg: cfg, baseURL: graphBaseURL, http: cfg.HTTPClient, dispatchSem: make(chan struct{}, maxConcurrentDispatch), readSem: make(chan struct{}, maxConcurrentReads)}, nil
}

func (a *adapter) Connect(ctx context.Context, deps core.AdapterDeps) error {
Expand Down Expand Up @@ -288,6 +299,19 @@ func (a *adapter) handleVerify(w http.ResponseWriter, r *http.Request) {
}

func (a *adapter) handleWebhook(dispatchCtx context.Context, w http.ResponseWriter, r *http.Request, deps core.AdapterDeps) {
// Bound concurrent inbound reads before buffering the body: a flood of large
// bodies is shed with 503 (Meta retries) rather than buffering unbounded
// memory. readSem is allocated once in newAdapter, so this direct read races
// nothing. Separate from dispatchSem, which bounds post-ack dispatch.
select {
case a.readSem <- struct{}{}:
default:
a.log().Warn("whatsapp: inbound concurrency limit reached; shedding with 503", "limit", maxConcurrentReads)
w.WriteHeader(http.StatusServiceUnavailable)
return
}
defer func() { <-a.readSem }()

body, err := io.ReadAll(http.MaxBytesReader(w, r.Body, maxRequestBytes))
if err != nil {
w.WriteHeader(http.StatusBadRequest)
Expand Down Expand Up @@ -500,8 +524,10 @@ func (a *adapter) ResolveAttachmentURL(ctx context.Context, att core.Attachment)
return "", nil
}

url := fmt.Sprintf("%s/%s/%s", a.baseURL, a.cfg.GraphVersion, media.ID)
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
// PathEscape the wire-derived media ID so it cannot alter the request path;
// it stays confined to the pinned graph.facebook.com host either way.
endpoint := fmt.Sprintf("%s/%s/%s", a.baseURL, a.cfg.GraphVersion, url.PathEscape(media.ID))
req, err := http.NewRequestWithContext(ctx, http.MethodGet, endpoint, nil)
if err != nil {
return "", err
}
Expand Down
51 changes: 50 additions & 1 deletion internal/whatsapp/cloud/cloud_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ func validConfig() Config {
func testAdapter() *adapter {
cfg := validConfig()
cfg.Path = defaultPath
return &adapter{cfg: cfg, baseURL: graphBaseURL, http: http.DefaultClient, dispatchSem: make(chan struct{}, maxConcurrentDispatch)}
return &adapter{cfg: cfg, baseURL: graphBaseURL, http: http.DefaultClient, dispatchSem: make(chan struct{}, maxConcurrentDispatch), readSem: make(chan struct{}, maxConcurrentReads)}
}

// sign returns the X-Hub-Signature-256 header value for body under secret.
Expand Down Expand Up @@ -484,6 +484,32 @@ func TestHandleWebhook_SaturationReturns503(t *testing.T) {
a.drainDispatch(ctx)
}

// readSpy counts Read calls so a test can prove a shed request's body was
// never touched.
type readSpy struct{ reads int }

func (s *readSpy) Read([]byte) (int, error) { s.reads++; return 0, io.EOF }

// A saturated read semaphore must answer 503 before the (up to maxRequestBytes)
// body is buffered, so a flood of large POSTs cannot exhaust memory ahead of the
// signature check. Mirrors the github/gitlab sibling test.
func TestHandleWebhook_ReadSaturationReturns503(t *testing.T) {
a := testAdapter()
a.readSem = make(chan struct{}, 1)
a.readSem <- struct{}{} // occupy the only inbound slot

var got []*core.Message
spy := &readSpy{}
r := httptest.NewRequest(http.MethodPost, "/webhook", spy)
r.Header.Set(signatureHeader, sign(a.cfg.AppSecret, []byte(textWebhook)))
w := httptest.NewRecorder()
a.handleWebhook(context.Background(), w, r, captureDeps(&got, nil))

asserts.Equal(t, w.Code, http.StatusServiceUnavailable, "saturated read gate returns 503")
asserts.Equal(t, len(got), 0, "nothing dispatched")
asserts.Equal(t, spy.reads, 0, "shed request's body is never read")
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.

func TestHandleWebhook_StatusOnlyIgnored(t *testing.T) {
a := testAdapter()
var got []*core.Message
Expand Down Expand Up @@ -713,6 +739,29 @@ func TestResolveAttachmentURL(t *testing.T) {
asserts.Equal(t, gotAuth, "Bearer tok", "resolve should set the bearer token")
}

// A wire-derived media ID with URL-significant characters must be PathEscaped so
// it cannot alter the request path (defense-in-depth; it can't escape the pinned
// host regardless).
func TestResolveAttachmentURL_EscapesMediaID(t *testing.T) {
var gotPath string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
gotPath = r.URL.EscapedPath()
_, _ = io.WriteString(w, `{"url":"https://lookaside.fbsbx.com/media/blob"}`)
}))
defer srv.Close()

a := testAdapter()
a.baseURL = srv.URL
a.http = srv.Client()
a.cfg.GraphVersion = "v23.0"

_, err := a.ResolveAttachmentURL(context.Background(),
core.Attachment{ExtraData: &Media{ID: "a/b?c"}})

asserts.NoError(t, err, "resolve should succeed")
asserts.Equal(t, gotPath, "/v23.0/a%2Fb%3Fc", "media id is PathEscaped in the request path")
}

func TestResolveAttachmentURL_NoMediaID(t *testing.T) {
a := testAdapter()
a.http = nil // a no-op resolve must short-circuit before any HTTP call
Expand Down
Loading