check / check (push) Successful in 1m5s
`GET /api/v1/chats/{id}/messages?count=N` returns a chat's recent messages, oldest first (`count` 1 to 100, default 20), and `POST` to the same path sends a text message and returns it once the chat client has taken it. Both answer `404` for any id the chats list does not show, looked up in the same contact list. The message record is defined once, for the webhooks to reuse.
Disclosures: items are read five at a time, because one item can be many times its text and 100 at once could pass the 16 MiB read limit; text too long gets `413`, a contact who deleted the chat `409` (unverified for one still connecting); a query the server cannot read gets its own `400`.
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
|
|
}
|