Webhooks: post each new incoming message (closes #7)
check / check (push) Successful in 2m3s

Each message a contact sends in a direct chat is posted to each webhook registered on that chat: `POST`, `Content-Type: application/json`, body `{"chat_id":N,"message":{...}}` with the same record the messages endpoint returns. The event handler only queues the message; four goroutines post it, one attempt each with a 10-second timeout and no retries, and with 100 already waiting a message is dropped and logged, so a slow webhook never delays the bot's replies. Posting stops last, after the chat client.

Disclosures: redirects are not followed; a logged failure can name the webhook's host but never its path, query or credentials; posts still in flight at shutdown are abandoned and logged.

Model: opus-5-5
This commit was merged in pull request #19.
This commit is contained in:
2026-09-29 10:23:24 +02:00
parent 0b9121a806
commit fa2aa09769
8 changed files with 816 additions and 54 deletions
+4 -1
View File
@@ -1,6 +1,7 @@
// Package api is the bot's HTTP API, through which another program
// reads the bot's chats, sends messages in them and registers webhooks
// on them.
// on them. Deliveries posts to those webhooks each message a contact
// sends.
//
// Every request must carry the credential, as "Authorization: Bearer
// {credential}"; with no credential configured, every request is
@@ -9,6 +10,8 @@
// Handlers call the chat client on the request's own goroutine. Doing
// so from the chat client's event handler would wait forever, since
// that goroutine also delivers the responses (see simplex.EventHandler).
// For the same reason, the event handler hands messages to Deliveries
// without waiting.
package api
import (
+187
View File
@@ -0,0 +1,187 @@
package api
import (
"bytes"
"context"
"encoding/json"
"errors"
"log/slog"
"net/http"
"net/url"
"sync"
"time"
"sneak.berlin/go/simplexcalc/internal/simplex"
)
const (
// queueLength is how many messages can wait to be posted, and
// deliveryGoroutines how many goroutines post them, each one message
// at a time.
queueLength = 100
deliveryGoroutines = 4
// deliveryTimeout bounds one POST to a webhook.
deliveryTimeout = 10 * time.Second
)
// delivery is the body of the POST to each webhook of a chat: a message a
// contact sent in it, as GET /api/v1/chats/{id}/messages shows it.
type delivery struct {
ChatID int64 `json:"chat_id"`
Message message `json:"message"`
}
// Deliveries posts each message a contact sends to each webhook of its
// chat, once, with no retry. Add hands it a message without waiting: the
// message joins a queue served by deliveryGoroutines goroutines, or is
// dropped if the queue is full. What goes wrong is logged with ids only:
// never a webhook's URL, which can hold a secret, nor a message's text.
type Deliveries struct {
log *slog.Logger
webhooks *Webhooks
client *http.Client
queue chan delivery
// stop ends the goroutines' context, abandoning the POSTs in
// progress; wg waits for the goroutines to return.
stop context.CancelFunc
wg sync.WaitGroup
}
// StartDeliveries starts posting the messages Add is given, until ctx
// ends or Stop is called.
func StartDeliveries(
ctx context.Context, log *slog.Logger, webhooks *Webhooks,
) *Deliveries {
ctx, stop := context.WithCancel(ctx)
d := &Deliveries{
log: log,
webhooks: webhooks,
client: &http.Client{
Timeout: deliveryTimeout,
// A redirect is not followed, and so counts as a status
// other than 2xx: following it could turn the POST into a
// GET without the message.
CheckRedirect: func(*http.Request, []*http.Request) error {
return http.ErrUseLastResponse
},
},
queue: make(chan delivery, queueLength),
stop: stop,
}
for range deliveryGoroutines {
d.wg.Go(func() { d.serve(ctx) })
}
return d
}
// Add queues the message in item for the webhooks of its chat, if it is a
// message a contact sent in a direct chat: one that GET
// /api/v1/chats/{id}/messages shows as received. The chat client's event
// handler calls it, so it never waits: if the queue is full, the message
// is dropped, and the log says so.
func (d *Deliveries) Add(item simplex.AChatItem) {
contact := item.ChatInfo.Contact
m, ok := newMessage(item.ChatItem)
if !ok || m.Direction != "received" || item.ChatInfo.Type != "direct" ||
contact == nil {
return
}
select {
case d.queue <- delivery{ChatID: contact.ContactID, Message: m}:
default:
d.log.Warn("dropping a message for the webhooks: the queue is full",
"chat_id", contact.ContactID, "message_id", m.ID)
}
}
// Stop stops posting, and returns once the goroutines have: POSTs in
// progress are abandoned, and messages still queued are not posted.
func (d *Deliveries) Stop() {
d.stop()
d.wg.Wait()
}
// serve posts queued messages, one at a time, until ctx ends.
func (d *Deliveries) serve(ctx context.Context) {
for {
select {
case <-ctx.Done():
return
case m := <-d.queue:
d.deliver(ctx, m)
}
}
}
// deliver posts m to each webhook of its chat, and logs each POST that
// fails or is answered with a status other than 2xx. The webhooks are
// read here rather than in Add, as reading them waits while the webhooks
// file is written.
func (d *Deliveries) deliver(ctx context.Context, m delivery) {
body, err := json.Marshal(m)
if err != nil {
d.log.Error("encoding a message for the webhooks",
"message_id", m.Message.ID, "error", err)
return
}
for _, hook := range d.webhooks.list(m.ChatID) {
// Stopped: the chat's other webhooks do not get the message.
if ctx.Err() != nil {
return
}
status, err := d.post(ctx, hook.URL, body)
switch {
case err != nil:
d.log.Warn("posting a message to a webhook",
"webhook_id", hook.ID, "message_id", m.Message.ID, "error", err)
case status < http.StatusOK || status >= http.StatusMultipleChoices:
d.log.Warn("posting a message to a webhook",
"webhook_id", hook.ID, "message_id", m.Message.ID, "status", status)
}
}
}
// post sends body to hookURL, once, and returns the status of the answer.
// Its error, unlike net/http's, never holds the URL.
func (d *Deliveries) post(
ctx context.Context, hookURL string, body []byte,
) (int, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodPost, hookURL,
bytes.NewReader(body))
if err != nil {
return 0, withoutURL(err)
}
req.Header.Set("Content-Type", "application/json")
resp, err := d.client.Do(req)
if err != nil {
return 0, withoutURL(err)
}
_ = resp.Body.Close()
return resp.StatusCode, nil
}
// withoutURL returns the cause of a *url.Error, whose own text holds the
// URL, and any other error as it is.
func withoutURL(err error) error {
var urlErr *url.Error
if errors.As(err, &urlErr) {
return urlErr.Err
}
return err
}
+328
View File
@@ -0,0 +1,328 @@
package api_test
import (
"bytes"
"encoding/json"
"io"
"log/slog"
"net/http"
"net/http/httptest"
"slices"
"strconv"
"strings"
"sync"
"testing"
"time"
"sneak.berlin/go/simplexcalc/internal/api"
"sneak.berlin/go/simplexcalc/internal/simplex"
)
// newChatItems is a newChatItems event, reduced to the fields the bot
// reads. It holds, in this order, a message in a group, a message the bot
// sent, an event the chat client records in a chat, and two messages a
// contact sent: a text and a picture with a caption.
const newChatItems = `{"type":"newChatItems","chatItems":[
{"chatInfo":{"type":"group","groupInfo":{"groupId":9}},
"chatItem":{"chatDir":{"type":"groupRcv"},"meta":{"itemId":5},
"content":{"type":"rcvMsgContent","msgContent":{"type":"text","text":"1 + 1"}}}},
{"chatInfo":{"type":"direct","contact":{"contactId":3}},
"chatItem":{"chatDir":{"type":"directSnd"},"meta":{"itemId":6},
"content":{"type":"sndMsgContent","msgContent":{"type":"text","text":"2"}}}},
{"chatInfo":{"type":"direct","contact":{"contactId":3}},
"chatItem":{"chatDir":{"type":"directRcv"},"meta":{"itemId":7},
"content":{"type":"rcvChatFeature","feature":"calls"}}},
{"chatInfo":{"type":"direct","contact":{"contactId":3}},
"chatItem":{"chatDir":{"type":"directRcv"},
"meta":{"itemId":9,"itemTs":"2026-09-29T03:14:34Z"},
"content":{"type":"rcvMsgContent","msgContent":{"type":"text","text":"2 + 2"}}}},
{"chatInfo":{"type":"direct","contact":{"contactId":3}},
"chatItem":{"chatDir":{"type":"directRcv"},
"meta":{"itemId":11,"itemTs":"2026-09-29T03:15:02Z"},
"content":{"type":"rcvMsgContent","msgContent":{"type":"image",
"text":"a picture","image":"data:image/jpg;base64,/9j/4AAQ"}}}}]}`
// textItem is the index in newChatItems of the text the contact sent.
const textItem = 3
// logBuffer is a log that the goroutines posting messages write to while
// the test reads it.
type logBuffer struct {
mu sync.Mutex
buf bytes.Buffer
}
func (b *logBuffer) Write(p []byte) (int, error) {
b.mu.Lock()
defer b.mu.Unlock()
return b.buf.Write(p)
}
func (b *logBuffer) String() string {
b.mu.Lock()
defer b.mu.Unlock()
return b.buf.String()
}
// chatItems returns the chat items in a newChatItems record.
func chatItems(t *testing.T, record string) []simplex.AChatItem {
t.Helper()
var r simplex.NewChatItems
err := json.Unmarshal([]byte(record), &r)
if err != nil {
t.Fatalf("decoding %s: %v", record, err)
}
return r.ChatItems
}
// receive returns the next of posts, and fails the test if none comes
// within 5 seconds.
func receive(t *testing.T, posts <-chan string) string {
t.Helper()
select {
case p := <-posts:
return p
case <-time.After(5 * time.Second):
t.Fatal("nothing more was posted")
return ""
}
}
// within runs f, and fails the test with complaint if f has not returned
// within 5 seconds.
func within(t *testing.T, complaint string, f func()) {
t.Helper()
done := make(chan struct{})
go func() {
defer close(done)
f()
}()
select {
case <-done:
case <-time.After(5 * time.Second):
t.Fatal(complaint)
}
}
// TestDeliveries: each message a contact sends is posted, as JSON with
// its chat's id, to each webhook of its chat and to no other chat's.
// Messages in groups, the bot's own messages and the events the chat
// client records in a chat are not posted.
func TestDeliveries(t *testing.T) {
t.Parallel()
posts := make(chan string, 10)
receiver := httptest.NewServer(http.HandlerFunc(
func(_ http.ResponseWriter, r *http.Request) {
body, _ := io.ReadAll(r.Body)
posts <- r.Method + " " + r.URL.Path + " " +
r.Header.Get("Content-Type") + " " + string(body)
}))
t.Cleanup(receiver.Close)
webhooks := readWebhooks(t, t.TempDir())
srv := webhookAPI(webhooks)
// Chat 4's webhook comes between chat 3's, so a message posted to
// every webhook in turn would reach it before chat 3's second.
register(t, srv, webhooksPath, receiver.URL+"/a", http.StatusCreated)
register(t, srv, "/api/v1/chats/4/webhooks", receiver.URL+"/c", http.StatusCreated)
register(t, srv, webhooksPath, receiver.URL+"/b", http.StatusCreated)
deliveries := api.StartDeliveries(t.Context(), slog.New(slog.DiscardHandler),
webhooks)
t.Cleanup(deliveries.Stop)
for _, item := range chatItems(t, newChatItems) {
deliveries.Add(item)
}
text := `{"chat_id":3,"message":{"id":9,"direction":"received",` +
`"type":"text","text":"2 + 2","time":"2026-09-29T03:14:34Z"}}`
picture := `{"chat_id":3,"message":{"id":11,"direction":"received",` +
`"type":"image","text":"a picture","time":"2026-09-29T03:15:02Z"}}`
want := []string{
"POST /a application/json " + text,
"POST /a application/json " + picture,
"POST /b application/json " + text,
"POST /b application/json " + picture,
}
got := make([]string, 0, len(want)+1)
for range want {
got = append(got, receive(t, posts))
}
deliveries.Stop()
select {
case p := <-posts:
got = append(got, p)
default:
}
slices.Sort(want)
slices.Sort(got)
if !slices.Equal(got, want) {
t.Errorf("posted:\n%s\nwant:\n%s",
strings.Join(got, "\n"), strings.Join(want, "\n"))
}
}
// TestDeliveriesQueueFull: while a webhook holds a POST from each of the
// goroutines, Add queues as many messages as the queue takes, then drops
// the next rather than wait, logging it by its chat's id and its own,
// never its text. Stop abandons the POSTs in progress rather than wait
// for them.
func TestDeliveriesQueueFull(t *testing.T) {
t.Parallel()
held := make(chan string, api.DeliveryGoroutines)
receiver := httptest.NewServer(http.HandlerFunc(
func(_ http.ResponseWriter, r *http.Request) {
// Read to the end: only then does the server notice when the
// bot abandons the POST, and end r's context.
body, _ := io.ReadAll(r.Body)
held <- string(body)
<-r.Context().Done() // never answers
}))
t.Cleanup(receiver.Close)
webhooks := readWebhooks(t, t.TempDir())
register(t, webhookAPI(webhooks), webhooksPath, receiver.URL, http.StatusCreated)
var logged logBuffer
deliveries := api.StartDeliveries(t.Context(),
slog.New(slog.NewJSONHandler(&logged, nil)), webhooks)
t.Cleanup(deliveries.Stop)
// A message for each goroutine, whose POST the webhook holds, then as
// many as the queue takes, then one more.
items := make([]simplex.AChatItem, api.DeliveryGoroutines+api.QueueLength+1)
for i := range items {
items[i] = chatItems(t, newChatItems)[textItem]
items[i].ChatItem.Meta.ItemID = int64(i + 1)
}
for _, item := range items[:api.DeliveryGoroutines] {
deliveries.Add(item)
receive(t, held)
}
within(t, "Add waited for room in the queue", func() {
for _, item := range items[api.DeliveryGoroutines:] {
deliveries.Add(item)
}
})
within(t, "Stop waited for the webhook to answer", deliveries.Stop)
log := logged.String()
dropped := `"level":"WARN","msg":"dropping a message for the webhooks: ` +
`the queue is full","chat_id":3,"message_id":`
if strings.Count(log, dropped) != 1 ||
!strings.Contains(log, dropped+strconv.Itoa(len(items))+"}") {
t.Errorf("log = %s, want message %d dropped, and no other", log, len(items))
}
if strings.Contains(log, "2 + 2") {
t.Errorf("log = %s, holding a message's text", log)
}
}
// TestDeliveryFailures: a POST that fails, or is answered with a status
// other than 2xx, is logged with the webhook's id, the message's and the
// reason or the status, and is not made again. A redirect is not
// followed. The log holds neither the webhook's URL nor the message's
// text.
func TestDeliveryFailures(t *testing.T) {
t.Parallel()
requested := make(chan string, 10)
receiver := httptest.NewServer(http.HandlerFunc(
func(w http.ResponseWriter, r *http.Request) {
requested <- r.URL.Path
switch r.URL.Path {
case "/refuses":
w.WriteHeader(http.StatusServiceUnavailable)
case "/redirects":
http.Redirect(w, r, "/elsewhere", http.StatusFound)
case "/hangs-up":
conn, _, err := http.NewResponseController(w).Hijack()
if err == nil {
_ = conn.Close()
}
}
}))
t.Cleanup(receiver.Close)
webhooks := readWebhooks(t, t.TempDir())
paths := []string{"/hangs-up", "/redirects", "/refuses"}
ids := map[string]string{}
// The query stands for a secret that a webhook's URL can hold.
for _, path := range paths {
ids[path] = register(t, webhookAPI(webhooks), webhooksPath,
receiver.URL+path+"?key=s3cret", http.StatusCreated).ID
}
var logged logBuffer
deliveries := api.StartDeliveries(t.Context(),
slog.New(slog.NewJSONHandler(&logged, nil)), webhooks)
t.Cleanup(deliveries.Stop)
deliveries.Add(chatItems(t, newChatItems)[textItem])
failed := `"level":"WARN","msg":"posting a message to a webhook","webhook_id":"`
within(t, "the three failures were not logged", func() {
for strings.Count(logged.String(), failed) < len(paths) {
time.Sleep(10 * time.Millisecond)
}
})
log := logged.String()
for _, want := range []string{
failed + ids["/hangs-up"] + `","message_id":9,"error":"`,
failed + ids["/redirects"] + `","message_id":9,"status":302}`,
failed + ids["/refuses"] + `","message_id":9,"status":503}`,
} {
if !strings.Contains(log, want) {
t.Errorf("log = %s, want it to hold %s", log, want)
}
}
for _, secret := range []string{receiver.URL, "s3cret", "2 + 2"} {
if strings.Contains(log, secret) {
t.Errorf("log = %s, holding %s", log, secret)
}
}
// Each webhook was asked once, before its failure was logged, and
// the redirect's target never.
got := []string{<-requested, <-requested, <-requested}
slices.Sort(got)
if !slices.Equal(got, paths) || len(requested) != 0 {
t.Errorf("requested %v and %d more, want %v once each",
got, len(requested), paths)
}
}
+9
View File
@@ -0,0 +1,9 @@
package api
// DeliveryGoroutines and QueueLength are deliveryGoroutines and
// queueLength, exported for the external test package, which fills the
// queue.
const (
DeliveryGoroutines = deliveryGoroutines
QueueLength = queueLength
)