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

Each message a contact sends in a direct chat is posted to each webhook
on that chat, as JSON holding `chat_id` and the message record of the
messages endpoint. The event handler only queues it: `api.Deliveries`
serves a queue of 100 with four goroutines, reads the chat's webhooks
there, and makes one POST per webhook, with a 10-second timeout, no
retry and no redirect followed. A full queue or a failed post is logged
by ids, never by URL or text. `bot.Run` stops it last, once the chat
client has exited. Tests cover the payload and its routing, a full
queue, failures, and a webhook that never answers; the README documents
the payload and what delivery promises.

Model: opus-5-5
This commit is contained in:
2026-09-29 07:19:56 +00:00
parent 397fc95149
commit 7170576a8e
8 changed files with 816 additions and 54 deletions
+56 -4
View File
@@ -226,7 +226,9 @@ Nothing is sent, and the answer is:
Registers the `url` in the body as a webhook on the chat `id`, and Registers the `url` in the body as a webhook on the chat `id`, and
answers `201` with the webhook. If the chat already has a webhook with answers `201` with the webhook. If the chat already has a webhook with
the same `url`, character for character, the answer is `200` with that the same `url`, character for character, the answer is `200` with that
webhook instead. webhook instead. Each new message the contact sends in the chat is then
posted to the webhook, as
[Messages posted to webhooks](#messages-posted-to-webhooks) describes.
```sh ```sh
curl -H "Authorization: Bearer $TOKEN" -H 'Content-Type: application/json' \ curl -H "Authorization: Bearer $TOKEN" -H 'Content-Type: application/json' \
@@ -252,8 +254,7 @@ The bot keeps the webhooks in `webhooks.json` in `DATA_DIR`, with mode
0600, and reads that file at startup, so they survive a restart. Without 0600, and reads that file at startup, so they survive a restart. Without
the file there are none; a file that cannot be read, or that holds the file there are none; a file that cannot be read, or that holds
anything but webhooks as the bot writes them, aborts startup, as anything but webhooks as the bot writes them, aborts startup, as
configuration that cannot be parsed does. The bot does not post anything configuration that cannot be parsed does.
to a webhook yet.
Nothing is registered, and the answer is: Nothing is registered, and the answer is:
@@ -303,6 +304,51 @@ Nothing is removed, and the answer is:
webhook `webhook_id`; webhook `webhook_id`;
- `500` if `webhooks.json` cannot be written. - `500` if `webhooks.json` cannot be written.
### Messages posted to webhooks
Each new message a contact sends in a chat is posted to each of the
chat's webhooks: a `POST` to its `url`, with
`Content-Type: application/json`, whose body is the chat's `id` as
`chat_id` and the message as `GET /api/v1/chats/{id}/messages` shows it:
```json
{
"chat_id": 3,
"message": {
"id": 9,
"direction": "received",
"type": "text",
"text": "2 + 2",
"time": "2026-09-29T03:14:34Z"
}
}
```
Every kind of message is posted, not only arithmetic, so `type` can be
`image`, `file` or another kind, with the caption as `text`. Messages
the bot sends, its replies and those sent through the API, are not
posted.
A message is posted to each webhook at most once: there are no retries.
A webhook misses it if the bot cannot reach it, if it does not answer
within 10 seconds, or if it answers with a status other than 2xx; a
redirect is not followed, and counts as such a status. So that no
webhook holds up the bot's replies, messages wait in a queue, and the
bot posts four at a time; a message that arrives while 100 wait is not
posted. Posting goes on while the bot stops; once its chat client has
exited, posts in progress are abandoned, and messages still waiting are
not posted.
The bot logs two warnings, which never hold a webhook's URL, as it can
carry a secret, nor a message's text:
- `dropping a message for the webhooks: the queue is full`, with
`chat_id` and `message_id`, for a message that found the queue full;
- `posting a message to a webhook`, with `webhook_id`, `message_id` and
either the `status` the webhook answered with or the `error` that
ended the post, for each post that got no 2xx answer, abandoned posts
included. The error can name the webhook's host.
## Entrypoints ## Entrypoints
This repo adheres to the This repo adheres to the
@@ -405,7 +451,13 @@ container.
each rewrites the file whole: the bot writes a temporary file beside each rewrites the file whole: the bot writes a temporary file beside
it, named `webhooks.json.` and digits, and renames it over the old it, named `webhooks.json.` and digits, and renames it over the old
one. A crash leaves the old file or the new one, never part of one, one. A crash leaves the old file or the new one, never part of one,
and at worst a stray temporary file, which the bot ignores. and at worst a stray temporary file, which the bot ignores. Messages
are posted to the webhooks by four goroutines serving a queue. The
goroutine that delivers the chat client's events, and also its
answers, only puts each message in the queue, never waiting: not for a
webhook, and not for the list of webhooks, which waits while the file
is written. Posting goes on while the bot stops, until the chat client
has exited.
- **Replies**: for each text message a contact sends in a direct chat, - **Replies**: for each text message a contact sends in a direct chat,
the bot sends back the result, as a reply quoting the message. Group the bot sends back the result, as a reply quoting the message. Group
messages, files and the bot's own messages are ignored. messages, files and the bot's own messages are ignored.
+1
View File
@@ -27,6 +27,7 @@ with no deprecation warning.
# Completed Steps # Completed Steps
- 2026-09-29 Each new incoming message is posted to its chat's webhooks
- 2026-09-29 `POST`, `GET` and `DELETE` on a chat's webhooks, under - 2026-09-29 `POST`, `GET` and `DELETE` on a chat's webhooks, under
`/api/v1/chats/{id}/webhooks`, kept in `$DATA_DIR/webhooks.json` `/api/v1/chats/{id}/webhooks`, kept in `$DATA_DIR/webhooks.json`
- 2026-09-29 `GET` and `POST /api/v1/chats/{id}/messages`: a chat's - 2026-09-29 `GET` and `POST /api/v1/chats/{id}/messages`: a chat's
+4 -1
View File
@@ -1,6 +1,7 @@
// Package api is the bot's HTTP API, through which another program // Package api is the bot's HTTP API, through which another program
// reads the bot's chats, sends messages in them and registers webhooks // 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 // Every request must carry the credential, as "Authorization: Bearer
// {credential}"; with no credential configured, every request is // {credential}"; with no credential configured, every request is
@@ -9,6 +10,8 @@
// Handlers call the chat client on the request's own goroutine. Doing // Handlers call the chat client on the request's own goroutine. Doing
// so from the chat client's event handler would wait forever, since // so from the chat client's event handler would wait forever, since
// that goroutine also delivers the responses (see simplex.EventHandler). // that goroutine also delivers the responses (see simplex.EventHandler).
// For the same reason, the event handler hands messages to Deliveries
// without waiting.
package api package api
import ( 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
)
+20 -9
View File
@@ -56,10 +56,10 @@ var errExited = errors.New("simplex-chat exited")
// Run reads the webhooks kept in cfg.DataDir, starts the chat client with // Run reads the webhooks kept in cfg.DataDir, starts the chat client with
// its database there, serving its API on localhost at chatPort, connects // its database there, serving its API on localhost at chatPort, connects
// to it, sets up the bot's address, then answers messages and serves the // to it, sets up the bot's address, then answers messages, posts them to
// bot's API until ctx is cancelled — which is a clean stop and returns // the webhooks and serves the bot's API until ctx is cancelled — which is
// nil — or until the chat client, the connection to it or the API's // a clean stop and returns nil — or until the chat client, the connection
// listener fails, which returns the error. // to it or the API's listener fails, which returns the error.
func Run( func Run(
ctx context.Context, log *slog.Logger, cfg *config.Config, chatPort int, ctx context.Context, log *slog.Logger, cfg *config.Config, chatPort int,
) error { ) error {
@@ -75,6 +75,12 @@ func Run(
return err return err
} }
// Posts the messages the event handler hands it to the webhooks. Its
// stop, deferred first, runs last, once the chat client has exited
// and can send no more messages; cancelling ctx does not reach it.
deliveries := api.StartDeliveries(context.WithoutCancel(ctx), log, webhooks)
defer deliveries.Stop()
// Cancelling this stops the chat client; the deferred wait makes // Cancelling this stops the chat client; the deferred wait makes
// Run return only once it has exited, whatever path Run takes. // Run return only once it has exited, whatever path Run takes.
// Cancelling ctx does not reach it, so that the API, stopped first, // Cancelling ctx does not reach it, so that the API, stopped first,
@@ -94,7 +100,7 @@ func Run(
<-cli.Done() <-cli.Done()
}() }()
client, err := connect(ctx, log, cli, chatPort) client, err := connect(ctx, log, cli, chatPort, handle(log, deliveries))
if err != nil { if err != nil {
return err return err
} }
@@ -155,9 +161,10 @@ func stopAPI(ctx context.Context, log *slog.Logger, srv *http.Server) {
} }
// connect waits for the chat client to open its API on chatPort and // connect waits for the chat client to open its API on chatPort and
// connects to it. // connects to it, passing its events to onEvent.
func connect( func connect(
ctx context.Context, log *slog.Logger, cli *simplex.CLI, chatPort int, ctx context.Context, log *slog.Logger, cli *simplex.CLI, chatPort int,
onEvent simplex.EventHandler,
) (*simplex.Client, error) { ) (*simplex.Client, error) {
ctx, cancel := context.WithTimeout(ctx, connectTimeout) ctx, cancel := context.WithTimeout(ctx, connectTimeout)
defer cancel() defer cancel()
@@ -165,7 +172,7 @@ func connect(
url := "ws://127.0.0.1:" + strconv.Itoa(chatPort) url := "ws://127.0.0.1:" + strconv.Itoa(chatPort)
for { for {
client, err := simplex.Dial(ctx, url, log, handle(log)) client, err := simplex.Dial(ctx, url, log, onEvent)
if err == nil { if err == nil {
return client, nil return client, nil
} }
@@ -228,8 +235,9 @@ func setUp(
return user, nil return user, nil
} }
// handle answers each text message a contact sends. // handle answers each text message a contact sends, and hands every
func handle(log *slog.Logger) simplex.EventHandler { // message a contact sends to deliveries, for the webhooks of its chat.
func handle(log *slog.Logger, deliveries *api.Deliveries) simplex.EventHandler {
return func(c *simplex.Client, ev simplex.Event) { return func(c *simplex.Client, ev simplex.Event) {
switch ev.Type { switch ev.Type {
case simplex.TypeNewChatItems: case simplex.TypeNewChatItems:
@@ -243,6 +251,9 @@ func handle(log *slog.Logger) simplex.EventHandler {
} }
for _, item := range r.ChatItems { for _, item := range r.ChatItems {
// Never waits, so no webhook holds up the replies.
deliveries.Add(item)
msg, ok := item.Message() msg, ok := item.Message()
if !ok { if !ok {
continue continue
+211 -40
View File
@@ -8,6 +8,7 @@ import (
"log/slog" "log/slog"
"net" "net"
"net/http" "net/http"
"net/http/httptest"
"net/netip" "net/netip"
"os" "os"
"path/filepath" "path/filepath"
@@ -35,11 +36,28 @@ const (
// long enough for the test to stop the bot meanwhile, and well // long enough for the test to stop the bot meanwhile, and well
// within the 5 seconds the API gets to finish its requests. // within the 5 seconds the API gets to finish its requests.
contactsDelay = time.Second contactsDelay = time.Second
// sent starts the name of the file the stand-in writes in the data
// directory for each message the bot sends, holding the command.
sent = "sent-"
) )
// twoMessages is the event the stand-in sends once the bot is set up: a
// contact sending the bot two messages. It is reduced to the fields the
// bot reads.
const twoMessages = `{"type":"newChatItems","chatItems":[
{"chatInfo":{"type":"direct","contact":{"contactId":3}},
"chatItem":{"chatDir":{"type":"directRcv"},
"meta":{"itemId":41,"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":42,"itemTs":"2026-09-29T03:14:35Z"},
"content":{"type":"rcvMsgContent","msgContent":{"type":"text","text":"3 * 3"}}}}]}`
// TestMain lets this test binary be the chat client as well: started // TestMain lets this test binary be the chat client as well: started
// under the chat client's name, as TestStopDuringRequest arranges, it is // under the chat client's name, as standInPath arranges, it is the
// the stand-in instead of running the tests. // stand-in instead of running the tests.
func TestMain(m *testing.M) { func TestMain(m *testing.M) {
if filepath.Base(os.Args[0]) == simplex.Binary { if filepath.Base(os.Args[0]) == simplex.Binary {
standIn() // never returns standIn() // never returns
@@ -58,13 +76,13 @@ func standIn() {
return os.Args[slices.Index(os.Args, name)+1] return os.Args[slices.Index(os.Args, name)+1]
} }
marker := filepath.Join(filepath.Dir(arg("--database")), asked) dir := filepath.Dir(arg("--database"))
srv := &http.Server{ srv := &http.Server{
Addr: "127.0.0.1:" + arg("--chat-server-port"), Addr: "127.0.0.1:" + arg("--chat-server-port"),
ReadHeaderTimeout: time.Second, ReadHeaderTimeout: time.Second,
Handler: http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { Handler: http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
answer(w, r, marker) answer(w, r, dir)
}), }),
} }
@@ -75,9 +93,11 @@ func standIn() {
} }
// answer answers the commands on one connection with records reduced to // answer answers the commands on one connection with records reduced to
// the fields the bot reads. Asked for the contacts, it first creates // the fields the bot reads. Asked for the contacts, it first creates the
// the file marker, then holds its answer for contactsDelay. // file asked in dir, then holds its answer for contactsDelay. Once the
func answer(w http.ResponseWriter, r *http.Request, marker string) { // bot is set up, it sends twoMessages. Each message the bot sends, it
// writes to a file in dir, named sent and the command's id.
func answer(w http.ResponseWriter, r *http.Request, dir string) {
conn, err := (&websocket.Upgrader{}).Upgrade(w, r, nil) conn, err := (&websocket.Upgrader{}).Upgrade(w, r, nil)
if err != nil { if err != nil {
return return
@@ -101,13 +121,20 @@ func answer(w http.ResponseWriter, r *http.Request, marker string) {
name, _, _ := strings.Cut(cmd["cmd"], " ") name, _, _ := strings.Cut(cmd["cmd"], " ")
if name == "/_send" {
_ = os.WriteFile(filepath.Join(dir, sent+cmd["corrId"]),
[]byte(cmd["cmd"]), 0o600)
continue
}
record, ok := records[name] record, ok := records[name]
if !ok { if !ok {
continue continue
} }
if name == "/_contacts" { if name == "/_contacts" {
_ = os.WriteFile(marker, nil, 0o600) _ = os.WriteFile(filepath.Join(dir, asked), nil, 0o600)
time.Sleep(contactsDelay) time.Sleep(contactsDelay)
} }
@@ -116,6 +143,11 @@ func answer(w http.ResponseWriter, r *http.Request, marker string) {
"corrId": cmd["corrId"], "corrId": cmd["corrId"],
"resp": json.RawMessage(record), "resp": json.RawMessage(record),
}) })
// The last command of the bot's set-up.
if name == "/_address_settings" {
_ = conn.WriteJSON(map[string]any{"resp": json.RawMessage(twoMessages)})
}
} }
} }
@@ -123,8 +155,129 @@ func answer(w http.ResponseWriter, r *http.Request, marker string) {
// when the bot is told to stop still gets the chat client's answer, // when the bot is told to stop still gets the chat client's answer,
// because the chat client is stopped only once the API has stopped. // because the chat client is stopped only once the API has stopped.
func TestStopDuringRequest(t *testing.T) { func TestStopDuringRequest(t *testing.T) {
// Run starts the chat client from PATH: put this test binary there t.Setenv("PATH", standInPath(t))
// under the chat client's name.
cfg := &config.Config{DataDir: t.TempDir(), Port: freePort(t), APIToken: credential}
stop, done := runBot(t, cfg)
// Stop the bot once the request below is waiting on the chat client.
go func() {
for {
_, err := os.Stat(filepath.Join(cfg.DataDir, asked))
if err == nil {
stop()
}
select {
case <-done:
return
case <-time.After(10 * time.Millisecond):
}
}
}()
status, body := getChats(t, cfg.Port, done)
want := `{"chats":[{"id":3,"display_name":"tester","contact_deleted":false}]}` + "\n"
if status != http.StatusOK || body != want {
t.Errorf("GET /api/v1/chats = %d %q, want 200 %q", status, body, want)
}
}
// TestSlowWebhook: a webhook that takes the bot's POSTs and never answers
// them holds up neither the bot's replies nor its stop.
func TestSlowWebhook(t *testing.T) {
t.Setenv("PATH", standInPath(t))
posts := make(chan string, 2)
receiver := httptest.NewServer(http.HandlerFunc(
func(_ http.ResponseWriter, r *http.Request) {
body, _ := io.ReadAll(r.Body)
posts <- string(body)
<-r.Context().Done() // until the bot abandons the POST
}))
t.Cleanup(receiver.Close)
cfg := &config.Config{DataDir: t.TempDir(), Port: freePort(t), APIToken: credential}
err := os.WriteFile(filepath.Join(cfg.DataDir, "webhooks.json"),
[]byte(`{"webhooks":[{"id":"00112233445566778899aabbccddeeff",`+
`"chat_id":3,"url":"`+receiver.URL+`"}]}`), 0o600)
if err != nil {
t.Fatal(err)
}
stop, done := runBot(t, cfg)
// The webhook gets both messages, and holds on to them...
got := []string{receive(t, posts, done), receive(t, posts, done)}
slices.Sort(got)
want := []string{
`{"chat_id":3,"message":{"id":41,"direction":"received",` +
`"type":"text","text":"2 + 2","time":"2026-09-29T03:14:34Z"}}`,
`{"chat_id":3,"message":{"id":42,"direction":"received",` +
`"type":"text","text":"3 * 3","time":"2026-09-29T03:14:35Z"}}`,
}
if !slices.Equal(got, want) {
t.Errorf("posted %q, want %q", got, want)
}
// ...while the bot answers both.
replies := []string{
`/_send @3 json [{"quotedItemId":41,` +
`"msgContent":{"type":"text","text":"4"},"mentions":{}}]`,
`/_send @3 json [{"quotedItemId":42,` +
`"msgContent":{"type":"text","text":"9"},"mentions":{}}]`,
}
deadline := time.Now().Add(5 * time.Second)
for !slices.Equal(sentMessages(t, cfg.DataDir), replies) {
if time.Now().After(deadline) {
t.Fatalf("sent %q, want %q", sentMessages(t, cfg.DataDir), replies)
}
time.Sleep(10 * time.Millisecond)
}
// Stopping abandons the POSTs, rather than give the webhook the 10
// seconds it has to answer.
stop()
select {
case <-done:
case <-time.After(5 * time.Second):
t.Fatal("Run did not return while a webhook held a POST")
}
}
// TestUnreadableWebhooks: a webhooks file that cannot be read stops Run
// before it starts the chat client.
func TestUnreadableWebhooks(t *testing.T) {
// No chat client on PATH: starting one would fail with another error.
t.Setenv("PATH", t.TempDir())
cfg := &config.Config{DataDir: t.TempDir(), Port: freePort(t)}
err := os.WriteFile(filepath.Join(cfg.DataDir, "webhooks.json"), []byte("{"), 0o600)
if err != nil {
t.Fatal(err)
}
err = bot.Run(t.Context(), slog.New(slog.DiscardHandler), cfg, freePort(t))
if err == nil || !strings.Contains(err.Error(), "webhooks.json") {
t.Errorf("Run = %v, want an error naming webhooks.json", err)
}
}
// standInPath returns a directory holding this test binary under the
// chat client's name. Run starts the chat client from PATH, so with PATH
// set to it, Run starts the stand-in.
func standInPath(t *testing.T) string {
t.Helper()
bin := t.TempDir() bin := t.TempDir()
exe, err := os.Executable() exe, err := os.Executable()
@@ -137,9 +290,15 @@ func TestStopDuringRequest(t *testing.T) {
t.Fatal(err) t.Fatal(err)
} }
t.Setenv("PATH", bin) return bin
}
// runBot runs the bot with cfg until stop is called or the test ends,
// and fails the test if Run returns an error. done is closed once Run
// has returned.
func runBot(t *testing.T, cfg *config.Config) (context.CancelFunc, <-chan struct{}) {
t.Helper()
cfg := &config.Config{DataDir: t.TempDir(), Port: freePort(t), APIToken: credential}
// Never bot.ChatPort: a real chat client may be listening there. // Never bot.ChatPort: a real chat client may be listening there.
chatPort := freePort(t) chatPort := freePort(t)
ctx, stop := context.WithCancel(t.Context()) ctx, stop := context.WithCancel(t.Context())
@@ -162,43 +321,55 @@ func TestStopDuringRequest(t *testing.T) {
} }
}) })
// Stop the bot once the request below is waiting on the chat client. return stop, done
go func() {
for ctx.Err() == nil {
_, err := os.Stat(filepath.Join(cfg.DataDir, asked))
if err == nil {
stop()
}
time.Sleep(10 * time.Millisecond)
}
}()
status, body := getChats(t, cfg.Port, done)
want := `{"chats":[{"id":3,"display_name":"tester","contact_deleted":false}]}` + "\n"
if status != http.StatusOK || body != want {
t.Errorf("GET /api/v1/chats = %d %q, want 200 %q", status, body, want)
}
} }
// TestUnreadableWebhooks: a webhooks file that cannot be read stops Run // receive returns the next of posts. It fails the test if none comes
// before it starts the chat client. // within 10 seconds, or if Run returns first, which closes done.
func TestUnreadableWebhooks(t *testing.T) { func receive(t *testing.T, posts <-chan string, done <-chan struct{}) string {
// No chat client on PATH: starting one would fail with another error. t.Helper()
t.Setenv("PATH", t.TempDir())
cfg := &config.Config{DataDir: t.TempDir(), Port: freePort(t)} select {
case p := <-posts:
return p
case <-done:
t.Fatal("Run returned before the webhook got the messages")
case <-time.After(10 * time.Second):
t.Fatal("the webhook did not get the messages")
}
err := os.WriteFile(filepath.Join(cfg.DataDir, "webhooks.json"), []byte("{"), 0o600) return ""
}
// sentMessages returns the commands of the messages the bot has sent, as
// the stand-in wrote them in dir, sorted.
func sentMessages(t *testing.T, dir string) []string {
t.Helper()
entries, err := os.ReadDir(dir)
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
err = bot.Run(t.Context(), slog.New(slog.DiscardHandler), cfg, freePort(t)) var cmds []string
if err == nil || !strings.Contains(err.Error(), "webhooks.json") {
t.Errorf("Run = %v, want an error naming webhooks.json", err) for _, entry := range entries {
if !strings.HasPrefix(entry.Name(), sent) {
continue
} }
//nolint:gosec // G304: the test's own file.
cmd, err := os.ReadFile(filepath.Join(dir, entry.Name()))
if err != nil {
t.Fatal(err)
}
cmds = append(cmds, string(cmd))
}
slices.Sort(cmds)
return cmds
} }
// freePort returns a TCP port that nothing listens on at the moment. // freePort returns a TCP port that nothing listens on at the moment.