From 7170576a8e79301dbae820b49d7d028b640e09c5 Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Tue, 29 Sep 2026 07:19:56 +0000 Subject: [PATCH] Webhooks: post each new incoming message (closes #7) 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 --- README.md | 60 +++++- docs/TODO.md | 1 + internal/api/api.go | 5 +- internal/api/deliveries.go | 187 ++++++++++++++++++ internal/api/deliveries_test.go | 328 ++++++++++++++++++++++++++++++++ internal/api/export_test.go | 9 + internal/bot/bot.go | 29 ++- internal/bot/run_test.go | 251 ++++++++++++++++++++---- 8 files changed, 816 insertions(+), 54 deletions(-) create mode 100644 internal/api/deliveries.go create mode 100644 internal/api/deliveries_test.go create mode 100644 internal/api/export_test.go diff --git a/README.md b/README.md index 331e690..ba6a764 100644 --- a/README.md +++ b/README.md @@ -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 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 -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 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 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 -configuration that cannot be parsed does. The bot does not post anything -to a webhook yet. +configuration that cannot be parsed does. Nothing is registered, and the answer is: @@ -303,6 +304,51 @@ Nothing is removed, and the answer is: webhook `webhook_id`; - `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 This repo adheres to the @@ -405,7 +451,13 @@ container. each rewrites the file whole: the bot writes a temporary file beside 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, - 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, the bot sends back the result, as a reply quoting the message. Group messages, files and the bot's own messages are ignored. diff --git a/docs/TODO.md b/docs/TODO.md index b2160f5..9a53b4f 100644 --- a/docs/TODO.md +++ b/docs/TODO.md @@ -27,6 +27,7 @@ with no deprecation warning. # 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 `/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 diff --git a/internal/api/api.go b/internal/api/api.go index 36d8eac..b4d6118 100644 --- a/internal/api/api.go +++ b/internal/api/api.go @@ -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 ( diff --git a/internal/api/deliveries.go b/internal/api/deliveries.go new file mode 100644 index 0000000..a037b54 --- /dev/null +++ b/internal/api/deliveries.go @@ -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 +} diff --git a/internal/api/deliveries_test.go b/internal/api/deliveries_test.go new file mode 100644 index 0000000..1a35513 --- /dev/null +++ b/internal/api/deliveries_test.go @@ -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) + } +} diff --git a/internal/api/export_test.go b/internal/api/export_test.go new file mode 100644 index 0000000..d718127 --- /dev/null +++ b/internal/api/export_test.go @@ -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 +) diff --git a/internal/bot/bot.go b/internal/bot/bot.go index 263ab51..6d343a1 100644 --- a/internal/bot/bot.go +++ b/internal/bot/bot.go @@ -56,10 +56,10 @@ var errExited = errors.New("simplex-chat exited") // Run reads the webhooks kept in cfg.DataDir, starts the chat client with // 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 -// bot's API until ctx is cancelled — which is a clean stop and returns -// nil — or until the chat client, the connection to it or the API's -// listener fails, which returns the error. +// to it, sets up the bot's address, then answers messages, posts them to +// the webhooks and serves the bot's API until ctx is cancelled — which is +// a clean stop and returns nil — or until the chat client, the connection +// to it or the API's listener fails, which returns the error. func Run( ctx context.Context, log *slog.Logger, cfg *config.Config, chatPort int, ) error { @@ -75,6 +75,12 @@ func Run( 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 // Run return only once it has exited, whatever path Run takes. // Cancelling ctx does not reach it, so that the API, stopped first, @@ -94,7 +100,7 @@ func Run( <-cli.Done() }() - client, err := connect(ctx, log, cli, chatPort) + client, err := connect(ctx, log, cli, chatPort, handle(log, deliveries)) if err != nil { 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 -// connects to it. +// connects to it, passing its events to onEvent. func connect( ctx context.Context, log *slog.Logger, cli *simplex.CLI, chatPort int, + onEvent simplex.EventHandler, ) (*simplex.Client, error) { ctx, cancel := context.WithTimeout(ctx, connectTimeout) defer cancel() @@ -165,7 +172,7 @@ func connect( url := "ws://127.0.0.1:" + strconv.Itoa(chatPort) for { - client, err := simplex.Dial(ctx, url, log, handle(log)) + client, err := simplex.Dial(ctx, url, log, onEvent) if err == nil { return client, nil } @@ -228,8 +235,9 @@ func setUp( return user, nil } -// handle answers each text message a contact sends. -func handle(log *slog.Logger) simplex.EventHandler { +// handle answers each text message a contact sends, and hands every +// 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) { switch ev.Type { case simplex.TypeNewChatItems: @@ -243,6 +251,9 @@ func handle(log *slog.Logger) simplex.EventHandler { } for _, item := range r.ChatItems { + // Never waits, so no webhook holds up the replies. + deliveries.Add(item) + msg, ok := item.Message() if !ok { continue diff --git a/internal/bot/run_test.go b/internal/bot/run_test.go index 76f78a8..67adae2 100644 --- a/internal/bot/run_test.go +++ b/internal/bot/run_test.go @@ -8,6 +8,7 @@ import ( "log/slog" "net" "net/http" + "net/http/httptest" "net/netip" "os" "path/filepath" @@ -35,11 +36,28 @@ const ( // long enough for the test to stop the bot meanwhile, and well // within the 5 seconds the API gets to finish its requests. 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 -// under the chat client's name, as TestStopDuringRequest arranges, it is -// the stand-in instead of running the tests. +// under the chat client's name, as standInPath arranges, it is the +// stand-in instead of running the tests. func TestMain(m *testing.M) { if filepath.Base(os.Args[0]) == simplex.Binary { standIn() // never returns @@ -58,13 +76,13 @@ func standIn() { 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{ Addr: "127.0.0.1:" + arg("--chat-server-port"), ReadHeaderTimeout: time.Second, 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 -// the fields the bot reads. Asked for the contacts, it first creates -// the file marker, then holds its answer for contactsDelay. -func answer(w http.ResponseWriter, r *http.Request, marker string) { +// the fields the bot reads. Asked for the contacts, it first creates the +// file asked in dir, then holds its answer for contactsDelay. Once the +// 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) if err != nil { return @@ -101,13 +121,20 @@ func answer(w http.ResponseWriter, r *http.Request, marker string) { 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] if !ok { continue } if name == "/_contacts" { - _ = os.WriteFile(marker, nil, 0o600) + _ = os.WriteFile(filepath.Join(dir, asked), nil, 0o600) time.Sleep(contactsDelay) } @@ -116,6 +143,11 @@ func answer(w http.ResponseWriter, r *http.Request, marker string) { "corrId": cmd["corrId"], "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, // because the chat client is stopped only once the API has stopped. func TestStopDuringRequest(t *testing.T) { - // Run starts the chat client from PATH: put this test binary there - // under the chat client's name. + t.Setenv("PATH", standInPath(t)) + + 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() exe, err := os.Executable() @@ -137,9 +290,15 @@ func TestStopDuringRequest(t *testing.T) { 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. chatPort := freePort(t) 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. - 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) - } + return stop, done } -// 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()) +// receive returns the next of posts. It fails the test if none comes +// within 10 seconds, or if Run returns first, which closes done. +func receive(t *testing.T, posts <-chan string, done <-chan struct{}) string { + t.Helper() - 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 { 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) + var cmds []string + + 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. -- 2.54.0