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 }