check / check (push) Successful in 1m21s
GET /api/v1/chats/{id}/messages returns the messages among a chat's
last count items (default 20, at most 100), oldest first. POST sends a
text, waits for the chat client's answer and returns 201 with the
message as sent. A chat the bot does not have is 404, a contact who
deleted the chat is 409, a text too long for one message is 413, and
any other failure is 500 with a chosen sentence.
The chat client spells out a message's formatting in its answer, up to
26 times the text's length, so 100 items in one answer can pass the
16 MiB read limit and end the connection. Items are read five at a
time.
Model: opus-5-5
465 lines
11 KiB
Go
465 lines
11 KiB
Go
// Package simplex runs the SimpleX Chat command-line client and talks
|
|
// to it over its WebSocket API.
|
|
//
|
|
// The protocol, documented in the simplex-chat repository under bots/:
|
|
// a command goes out as {"corrId": "...", "cmd": "..."}, and the client
|
|
// answers it with {"corrId": "...", "resp": {...}} carrying the same id.
|
|
// Everything it sends without a corrId is an event. The API has no
|
|
// authentication; the client binds it to localhost only.
|
|
package simplex
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"log/slog"
|
|
"slices"
|
|
"strconv"
|
|
"strings"
|
|
"sync"
|
|
|
|
"github.com/gorilla/websocket"
|
|
)
|
|
|
|
// maxMessageSize bounds one message from the chat client; a larger one
|
|
// ends the connection. The largest it sends are the pages of chat items
|
|
// ChatItems asks for, which chatItemsPage keeps under this.
|
|
const maxMessageSize = 16 << 20
|
|
|
|
// chatItemsPage is how many chat items ChatItems asks for at a time. An
|
|
// item repeats a message's text and spells out its formatting, which for
|
|
// text made of short mentions makes it 26 times as long as the text, and
|
|
// a contact's message can hold 64 KiB: one item can take 1.7 MiB. Five
|
|
// stay well under maxMessageSize.
|
|
const chatItemsPage = 5
|
|
|
|
var (
|
|
// ErrClosed is returned by Command once the connection has ended.
|
|
ErrClosed = errors.New("connection to the chat client closed")
|
|
|
|
// ErrNoContact is returned for a contact the user does not have.
|
|
ErrNoContact = errors.New("no such contact")
|
|
|
|
// ErrContactNotReady is returned for sending to a contact who cannot
|
|
// receive messages: one who has deleted their chat with the user, or
|
|
// who has not finished connecting.
|
|
ErrContactNotReady = errors.New("the contact cannot receive messages")
|
|
|
|
// ErrMessageTooLarge is returned for a message too large to send.
|
|
ErrMessageTooLarge = errors.New("the message is too large")
|
|
|
|
errUnexpected = errors.New("unexpected response")
|
|
errCommand = errors.New("command failed")
|
|
)
|
|
|
|
// EventHandler receives each event the chat client sends. Events are
|
|
// delivered one at a time, on the goroutine that also delivers command
|
|
// responses: a handler may Send, but must never wait on Command, whose
|
|
// response could then never arrive.
|
|
type EventHandler func(c *Client, ev Event)
|
|
|
|
// Client is a connection to the chat client's WebSocket API.
|
|
type Client struct {
|
|
conn *websocket.Conn
|
|
log *slog.Logger
|
|
onEvent EventHandler
|
|
|
|
// writeMu serialises writes: the connection allows one writer at
|
|
// a time, and Command and Send are called from different
|
|
// goroutines.
|
|
writeMu sync.Mutex
|
|
|
|
mu sync.Mutex
|
|
lastID uint64
|
|
waiting map[string]chan Event
|
|
|
|
done chan struct{}
|
|
err error // why the read loop ended; valid once done is closed
|
|
}
|
|
|
|
// Dial connects to the chat client's API at url and starts reading from
|
|
// it, passing every event to onEvent.
|
|
func Dial(
|
|
ctx context.Context, url string, log *slog.Logger, onEvent EventHandler,
|
|
) (*Client, error) {
|
|
conn, resp, err := websocket.DefaultDialer.DialContext(ctx, url, nil)
|
|
if resp != nil {
|
|
_ = resp.Body.Close()
|
|
}
|
|
|
|
if err != nil {
|
|
return nil, fmt.Errorf("connecting to %s: %w", url, err)
|
|
}
|
|
|
|
conn.SetReadLimit(maxMessageSize)
|
|
|
|
c := &Client{
|
|
conn: conn,
|
|
log: log,
|
|
onEvent: onEvent,
|
|
waiting: make(map[string]chan Event),
|
|
done: make(chan struct{}),
|
|
}
|
|
|
|
go c.read()
|
|
|
|
return c, nil
|
|
}
|
|
|
|
// Done is closed when the connection ends; Err then says why.
|
|
func (c *Client) Done() <-chan struct{} {
|
|
return c.done
|
|
}
|
|
|
|
// Err returns why the connection ended. Call it only after Done is
|
|
// closed.
|
|
func (c *Client) Err() error {
|
|
return c.err
|
|
}
|
|
|
|
// Close ends the connection.
|
|
func (c *Client) Close() error {
|
|
err := c.conn.Close()
|
|
if err != nil {
|
|
return fmt.Errorf("closing connection: %w", err)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// ActiveUser returns the chat client's active user profile.
|
|
func (c *Client) ActiveUser(ctx context.Context) (User, error) {
|
|
var r struct {
|
|
User User `json:"user"`
|
|
}
|
|
|
|
err := c.command(ctx, cmdShowActiveUser, TypeActiveUser, &r)
|
|
|
|
return r.User, err
|
|
}
|
|
|
|
// Address returns the user's long-term contact address, and false if
|
|
// the user has none.
|
|
func (c *Client) Address(ctx context.Context, userID int64) (ConnLink, bool, error) {
|
|
//nolint:tagliatelle // the chat client's wire format.
|
|
var r struct {
|
|
ContactLink struct {
|
|
ConnLinkContact ConnLink `json:"connLinkContact"`
|
|
} `json:"contactLink"`
|
|
}
|
|
|
|
err := c.command(ctx, cmdShowAddress(userID), TypeUserContactLink, &r)
|
|
|
|
var cerr *CommandError
|
|
if errors.As(err, &cerr) && cerr.Detail == "userContactLinkNotFound" {
|
|
return ConnLink{}, false, nil
|
|
}
|
|
|
|
if err != nil {
|
|
return ConnLink{}, false, err
|
|
}
|
|
|
|
return r.ContactLink.ConnLinkContact, true, nil
|
|
}
|
|
|
|
// CreateAddress creates the user's long-term contact address.
|
|
func (c *Client) CreateAddress(ctx context.Context, userID int64) (ConnLink, error) {
|
|
//nolint:tagliatelle // the chat client's wire format.
|
|
var r struct {
|
|
ConnLinkContact ConnLink `json:"connLinkContact"`
|
|
}
|
|
|
|
err := c.command(ctx, cmdCreateAddress(userID), TypeUserContactLinkCreated, &r)
|
|
|
|
return r.ConnLinkContact, err
|
|
}
|
|
|
|
// SetAddressSettings replaces the settings of the user's address.
|
|
func (c *Client) SetAddressSettings(
|
|
ctx context.Context, userID int64, s AddressSettings,
|
|
) error {
|
|
cmd, err := cmdSetAddressSettings(userID, s)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
return c.command(ctx, cmd, TypeUserContactLinkUpdated, nil)
|
|
}
|
|
|
|
// Contacts returns the user's contacts: everyone it has a direct chat
|
|
// with.
|
|
func (c *Client) Contacts(ctx context.Context, userID int64) ([]Contact, error) {
|
|
var r struct {
|
|
Contacts []Contact `json:"contacts"`
|
|
}
|
|
|
|
err := c.command(ctx, cmdListContacts(userID), TypeContactsList, &r)
|
|
|
|
return r.Contacts, err
|
|
}
|
|
|
|
// ChatItems returns the last count items of the chat with a contact,
|
|
// oldest first, or all of them if the chat has fewer.
|
|
func (c *Client) ChatItems(
|
|
ctx context.Context, contactID int64, count int,
|
|
) ([]ChatItem, error) {
|
|
var (
|
|
items []ChatItem
|
|
before int64 // 0 asks for the chat's last items
|
|
)
|
|
|
|
for len(items) < count {
|
|
page := min(count-len(items), chatItemsPage)
|
|
|
|
//nolint:tagliatelle // the chat client's wire format.
|
|
var r struct {
|
|
Chat struct {
|
|
ChatItems []ChatItem `json:"chatItems"`
|
|
} `json:"chat"`
|
|
}
|
|
|
|
err := c.command(ctx, cmdGetChat(contactID, before, page), TypeAPIChat, &r)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
items = slices.Concat(r.Chat.ChatItems, items)
|
|
|
|
if len(r.Chat.ChatItems) < page {
|
|
break
|
|
}
|
|
|
|
before = items[0].Meta.ItemID
|
|
}
|
|
|
|
return items, nil
|
|
}
|
|
|
|
// SendText sends a text message to a contact, as a reply to the message
|
|
// quotedItemID (0 for none). It does not wait for the chat client to
|
|
// accept it; a failure is logged when the client's answer arrives.
|
|
func (c *Client) SendText(contactID, quotedItemID int64, text string) error {
|
|
cmd, err := cmdSendText(contactID, quotedItemID, text)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
c.mu.Lock()
|
|
id := c.nextID()
|
|
c.mu.Unlock()
|
|
|
|
return c.write(id, cmd)
|
|
}
|
|
|
|
// SendMessage sends a text message to a contact and returns it as the
|
|
// chat client recorded it. Unlike SendText, it waits for the chat
|
|
// client's answer, so an EventHandler must never call it.
|
|
func (c *Client) SendMessage(
|
|
ctx context.Context, contactID int64, text string,
|
|
) (ChatItem, error) {
|
|
cmd, err := cmdSendText(contactID, 0, text)
|
|
if err != nil {
|
|
return ChatItem{}, err
|
|
}
|
|
|
|
var r NewChatItems
|
|
|
|
err = c.command(ctx, cmd, TypeNewChatItems, &r)
|
|
if err != nil {
|
|
return ChatItem{}, err
|
|
}
|
|
|
|
if len(r.ChatItems) != 1 {
|
|
return ChatItem{}, fmt.Errorf("%w to %q: %d chat items",
|
|
errUnexpected, cmdName(cmd), len(r.ChatItems))
|
|
}
|
|
|
|
return r.ChatItems[0].ChatItem, nil
|
|
}
|
|
|
|
// CommandError is a command the chat client refused. Type and Detail
|
|
// are the discriminators of its chatError record, such as "errorStore"
|
|
// and "userContactLinkNotFound".
|
|
type CommandError struct {
|
|
Type string
|
|
Detail string
|
|
}
|
|
|
|
func (e *CommandError) Error() string {
|
|
return fmt.Sprintf("%s: %s/%s", errCommand, e.Type, e.Detail)
|
|
}
|
|
|
|
func (e *CommandError) Unwrap() error {
|
|
return errCommand
|
|
}
|
|
|
|
// command sends cmd, waits for its response, and, if the response has
|
|
// type want, decodes it into out (unless out is nil).
|
|
func (c *Client) command(ctx context.Context, cmd, want string, out any) error {
|
|
ch := make(chan Event, 1)
|
|
|
|
c.mu.Lock()
|
|
id := c.nextID()
|
|
c.waiting[id] = ch
|
|
c.mu.Unlock()
|
|
|
|
defer func() {
|
|
c.mu.Lock()
|
|
delete(c.waiting, id)
|
|
c.mu.Unlock()
|
|
}()
|
|
|
|
err := c.write(id, cmd)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
var ev Event
|
|
|
|
select {
|
|
case ev = <-ch:
|
|
case <-c.done:
|
|
return fmt.Errorf("%w: %w", ErrClosed, c.err)
|
|
case <-ctx.Done():
|
|
return fmt.Errorf("waiting for a response to %q: %w", cmdName(cmd), ctx.Err())
|
|
}
|
|
|
|
switch ev.Type {
|
|
case want:
|
|
if out == nil {
|
|
return nil
|
|
}
|
|
|
|
return ev.Decode(out)
|
|
case TypeChatCmdError:
|
|
return commandError(ev)
|
|
default:
|
|
return fmt.Errorf("%w to %q: %s", errUnexpected, cmdName(cmd), ev.Type)
|
|
}
|
|
}
|
|
|
|
// nextID returns a fresh correlation id. The caller holds c.mu.
|
|
func (c *Client) nextID() string {
|
|
c.lastID++
|
|
|
|
return strconv.FormatUint(c.lastID, 10)
|
|
}
|
|
|
|
func (c *Client) write(id, cmd string) error {
|
|
c.writeMu.Lock()
|
|
defer c.writeMu.Unlock()
|
|
|
|
err := c.conn.WriteJSON(command{CorrID: id, Cmd: cmd})
|
|
if err != nil {
|
|
return fmt.Errorf("sending %q: %w", cmdName(cmd), err)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// read is the only reader of the connection. It runs until the
|
|
// connection fails or is closed.
|
|
func (c *Client) read() {
|
|
defer close(c.done)
|
|
|
|
for {
|
|
_, data, err := c.conn.ReadMessage()
|
|
if err != nil {
|
|
c.err = err
|
|
|
|
return
|
|
}
|
|
|
|
c.dispatch(data)
|
|
}
|
|
}
|
|
|
|
// dispatch routes one message: a response to whoever waits for it, an
|
|
// event to the handler. A message that does not parse is logged and
|
|
// skipped rather than ending the connection; the API documentation
|
|
// warns that records change between releases.
|
|
func (c *Client) dispatch(data []byte) {
|
|
var env envelope
|
|
|
|
err := json.Unmarshal(data, &env)
|
|
if err != nil {
|
|
c.log.Warn("undecodable message from the chat client", "error", err)
|
|
|
|
return
|
|
}
|
|
|
|
var head tagged
|
|
|
|
err = json.Unmarshal(env.Resp, &head)
|
|
if err != nil {
|
|
c.log.Warn("undecodable record from the chat client", "error", err)
|
|
|
|
return
|
|
}
|
|
|
|
ev := Event{Type: head.Type, raw: env.Resp}
|
|
|
|
if env.CorrID == "" {
|
|
c.onEvent(c, ev)
|
|
|
|
return
|
|
}
|
|
|
|
c.mu.Lock()
|
|
ch, ok := c.waiting[env.CorrID]
|
|
c.mu.Unlock()
|
|
|
|
if ok {
|
|
ch <- ev
|
|
|
|
return
|
|
}
|
|
|
|
// The response to a SendText: nobody waits for it, so a failure
|
|
// is reported here or nowhere.
|
|
if ev.Type == TypeChatCmdError {
|
|
c.log.Warn("sending a message failed", "error", commandError(ev))
|
|
}
|
|
}
|
|
|
|
// commandError returns the error in a chatCmdError record, marked with
|
|
// this package's error for the refusals that have one.
|
|
func commandError(ev Event) error {
|
|
var r cmdError
|
|
|
|
err := ev.Decode(&r)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
e := &CommandError{Type: r.ChatError.Type}
|
|
|
|
for _, detail := range []*tagged{
|
|
r.ChatError.ErrorType, r.ChatError.StoreError, r.ChatError.AgentError,
|
|
} {
|
|
if detail != nil {
|
|
e.Detail = detail.Type
|
|
}
|
|
}
|
|
|
|
switch e.Detail {
|
|
case "contactNotFound":
|
|
return fmt.Errorf("%w: %w", ErrNoContact, e)
|
|
case "contactNotReady":
|
|
return fmt.Errorf("%w: %w", ErrContactNotReady, e)
|
|
case "largeMsg":
|
|
return fmt.Errorf("%w: %w", ErrMessageTooLarge, e)
|
|
default:
|
|
return e
|
|
}
|
|
}
|
|
|
|
// cmdName is a command without its arguments, for error messages: the
|
|
// arguments of /_send are a message someone wrote.
|
|
func cmdName(cmd string) string {
|
|
name, _, _ := strings.Cut(cmd, " ")
|
|
|
|
return name
|
|
}
|