check / check (push) Successful in 54s
Remove the template's HTTP service, database and fx wiring. Add exact arithmetic on go/parser and go/constant, a client that runs simplex-chat as a child process and drives its WebSocket API, and the bot, which keeps an auto-accepting address and replies to each message. The image adds the checksum-pinned simplex-chat v7.0.2 on Ubuntu 22.04. Model: opus-5-5
360 lines
8.2 KiB
Go
360 lines
8.2 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"
|
|
"strconv"
|
|
"strings"
|
|
"sync"
|
|
|
|
"github.com/gorilla/websocket"
|
|
)
|
|
|
|
// maxMessageSize bounds one message from the chat client. The largest
|
|
// thing it sends is a record carrying a contact's profile picture, well
|
|
// under this.
|
|
const maxMessageSize = 16 << 20
|
|
|
|
var (
|
|
// ErrClosed is returned by Command once the connection has ended.
|
|
ErrClosed = errors.New("connection to the chat client closed")
|
|
|
|
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)
|
|
}
|
|
|
|
// 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)
|
|
}
|
|
|
|
// 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))
|
|
}
|
|
}
|
|
|
|
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
|
|
}
|
|
}
|
|
|
|
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
|
|
}
|